Skip to content
Satya Prakash Solanki

Part 1 of this series explained what TPC-H is and what its 22 queries exercise. Part 2 loaded the data into Iceberg and ran the queries through JDBC, capturing elapsed time, CPU time and rows processed per query.

That CSV tells you which queries are slow. It does not tell you why. The answer is in the query plan, and Trino is generous with plans once you know where to look. This part is about reading them. The workload is derived from TPC-H, and every number below is illustrative.

Three ways to ask for a plan

Trino offers several variants of EXPLAIN, and they answer different questions.

Statement Runs the query? Use it for
EXPLAIN No Logical shape: join order, filters, aggregations, with optimiser estimates
EXPLAIN (TYPE DISTRIBUTED) No How the plan is cut into fragments and where data moves between workers
EXPLAIN (TYPE IO) No Which tables and columns are read, and with which constraints
EXPLAIN ANALYZE Yes Actual rows, CPU, wall and blocked time, memory, per fragment and operator
EXPLAIN ANALYZE VERBOSE Yes Adds distributions across tasks and drivers, useful for skew

Two practical notes. First, EXPLAIN ANALYZE executes the statement, so it costs a full run and must never be pointed at an INSERT or DELETE you did not mean to execute. Second, the documentation warns that its statistics may not be entirely accurate for queries that finish quickly. At SF1, many TPC-H derived queries finish in a second or two, so I read plans at SF10 or above.

The anatomy of a distributed plan

A distributed plan is a tree of fragments. Each fragment runs as a stage, with one or more tasks across workers. The header of each fragment tells you how its work is distributed:

  • SOURCE fragments read tables. Their parallelism comes from splits.
  • HASH fragments receive rows partitioned by a hash of some keys. Partitioned joins and grouped aggregations live here.
  • SINGLE fragments run on one node. The final output and global ORDER BY with LIMIT end up here.
  • ROUND_ROBIN and broadcast-style outputs show up when data is spread evenly or replicated.

Fragments talk to each other through exchanges. A RemoteSource node is where a fragment reads the output of another fragment over the network. LocalExchange nodes repartition data between drivers inside one task. Every remote exchange is network traffic and a potential point of waiting, so the first thing I do with a slow query is count them and see what flows through each.

  1. 01Scan and filter (SOURCE)
  2. 02Repartition by key
  3. 03Join or aggregate (HASH)
  4. 04Gather
  5. 05Final sort and output (SINGLE)
Figure 1. The usual shape of a TPC-H derived query in Trino. The repartition step is where most network and skew problems show up.

Join distribution: replicated or partitioned

For every join, the cost-based optimiser chooses between two strategies:

  • Replicated (broadcast). The build side is sent in full to every worker that processes the probe side. Cheap when the build side is small, because the large probe side does not move. Every node must hold the whole build side in memory.
  • Partitioned. Both sides are redistributed by a hash of the join key. More network, but memory is spread across the cluster, so it scales to large build sides.

In a plan the join node carries distribution = REPLICATED or distribution = PARTITIONED. The choice is controlled by the join_distribution_type session property (AUTOMATIC by default, letting the optimiser decide from cost) and capped by join_max_broadcast_table_size. Join order is controlled separately by join_reordering_strategy.

In TPC-H derived queries, joins to nation, region and filtered part or supplier should almost always be replicated. Joins between lineitem and orders at any meaningful scale should be partitioned. When I see the opposite, the optimiser has been misled, and the usual reason is statistics.

Statistics and estimates

The optimiser’s choices are only as good as its inputs. EXPLAIN prints an Estimates line on each node, with rows, CPU, memory and network cost, and prints ? where it has no statistics (cost in EXPLAIN). A plan full of question marks means the optimiser is guessing.

SHOW STATS FOR lakehouse.tpch_sf100.orders;
-- if row counts or distinct values are missing:
ANALYZE lakehouse.tpch_sf100.orders;

For Iceberg tables loaded with CREATE TABLE AS SELECT, as in Part 2, the Iceberg connector can collect statistics on write, and ANALYZE collects them explicitly. I check SHOW STATS after every load because a benchmark run with statistics and one without are different experiments.

The most useful habit with EXPLAIN ANALYZE is comparing the estimate with the actual on each node. The documentation notes that estimated cost is printed alongside the runtime statistics. If an operator was estimated at a few thousand rows and produced hundreds of millions, everything above it was planned on a false premise. That one comparison is the most common explanation for a bad join order.

Dynamic filtering

Dynamic filtering lets Trino collect the join keys from the build side of a join and push them into the scan of the probe side while the query runs. For a star-shaped query, it can turn a full scan of lineitem into a scan of a small fraction of it.

In EXPLAIN you will see dynamicFilterAssignments on the join node and a dynamicFilters predicate on the scan it feeds. In EXPLAIN ANALYZE the scan’s input rows tell you whether the filter actually pruned anything. It is enabled by default and supported for inner and right joins and for semi-joins from IN conditions. When a filter is present but the scan still reads every row, check whether the build side completed before the probe scan started, and whether the connector could use the domain to skip files or row groups.

Skew and spilling

A plan can be correct and still slow because the work is uneven. With EXPLAIN ANALYZE, each fragment and operator reports input statistics as an average with a standard deviation across tasks. VERBOSE adds percentile distributions. A standard deviation close to or larger than the average means a few tasks are doing most of the work, and the stage takes as long as its slowest task.

TPC-H data is generated with uniform keys, so severe skew in a derived workload is itself a clue. It usually means the partitioning key is low-cardinality (grouping by n_name after an early repartition, for example) or a filter has collapsed the data onto a few values.

Memory shows up per operator as peak memory. If a hash aggregation or join build approaches the per-node memory limit under concurrency, the query will fail or, if spill to disk is enabled, start spilling. The documentation now describes spilling as legacy and recommends fault-tolerant execution with a task retry policy instead. Either way, spill or memory failure in a benchmark is a finding to explain.

Worked example: Q9, product type profit

Q9 joins six tables (part, supplier, lineitem, partsupp, orders, nation) and filters part on p_name LIKE '%green%'. It is usually one of the most expensive queries in the suite, and it rewards plan reading.

A healthy distributed plan has roughly this shape (abridged and illustrative; formatting differs between Trino versions):

Fragment 1 [HASH]
Aggregate(FINAL)[nation, o_year]
└─ RemoteSource[2]
Fragment 2 [HASH]
InnerJoin[o_orderkey = l_orderkey], distribution = PARTITIONED
├─ RemoteSource[3] -- lineitem joined with part, partsupp, supplier, nation
└─ RemoteSource[6] -- orders, repartitioned by orderkey
Fragment 3 [HASH]
InnerJoin[(l_partkey, l_suppkey) = (ps_partkey, ps_suppkey)], distribution = PARTITIONED
├─ InnerJoin[l_partkey = p_partkey], distribution = REPLICATED
│ dynamicFilterAssignments = {p_partkey -> #df_1}
│ ├─ ScanFilterProject[lineitem, dynamicFilters = {l_partkey = #df_1}]
│ └─ RemoteSource[4] -- part filtered on name, broadcast
└─ RemoteSource[5] -- partsupp

What to check, in order:

  1. Is the filtered part side replicated, and does it feed a dynamic filter into the lineitem scan? The LIKE '%green%' predicate keeps a small share of parts. If the dynamic filter works, the lineitem scan’s output rows in EXPLAIN ANALYZE drop sharply relative to its input.
  2. Is the lineitem to orders join partitioned? Replicating orders at SF100 would send a very large table to every worker.
  3. Are the estimates close? The selectivity of a LIKE with leading wildcard is hard to estimate. If the optimiser underestimates the filtered part rows badly, it may choose a poor join order. Comparing estimated and actual rows on the part scan shows this immediately.
  4. Where is the time? In a healthy run the lineitem scan and the two large joins dominate CPU. If the final aggregation or the nation join shows significant time, something upstream produced far more rows than it should.

Worked example: Q18, large volume customer

Q18 in the specification finds orders whose total quantity exceeds a threshold:

SELECT c_name, c_custkey, o_orderkey, o_orderdate, o_totalprice, sum(l_quantity)
FROM customer, orders, lineitem
WHERE o_orderkey IN (
SELECT l_orderkey FROM lineitem
GROUP BY l_orderkey HAVING sum(l_quantity) > 300)
AND c_custkey = o_custkey
AND o_orderkey = l_orderkey
GROUP BY c_name, c_custkey, o_orderkey, o_orderdate, o_totalprice
ORDER BY o_totalprice DESC, o_orderdate
LIMIT 100;

Its plan shape is quite different from Q9. The IN subquery becomes a semi-join, and the subquery itself is a grouped aggregation over all of lineitem by l_orderkey. That aggregation has one group per order, which at SF100 is a very large hash table, partitioned across workers.

What to check:

  1. Peak memory on the subquery aggregation. This is usually the largest memory consumer in the query and the first to fail or spill under concurrency.
  2. Whether the semi-join result feeds a dynamic filter. Only a small number of orders pass the HAVING condition. If that set is pushed into the scans of orders and the outer lineitem, the rest of the query is cheap. If not, the outer joins process far more rows than they need.
  3. Two scans of lineitem. One for the subquery and one for the outer join. That is expected; the question is whether the second is pruned.
  4. The final ORDER BY ... LIMIT. It should appear as a partial top-N on workers and a final top-N in a single fragment, not as a full sort.

Note that simplified variants of Q18 circulate widely, including some with a different HAVING condition. They have different plans, which is one more reason to pin the exact query text in a benchmark.

A symptom-to-cause checklist

Symptom in the plan Likely cause First action
? in estimates, odd join order Missing statistics SHOW STATS, then ANALYZE
Large table with distribution = REPLICATED Underestimated size Fix statistics; test join_distribution_type per session
Dynamic filter assigned but scan input unchanged Filter arrived late or could not prune Check build side timing and connector pruning
Std.dev of input close to or above average Skew on partition key Inspect keys; check repartition placement
Peak memory near per-node limit Large build or aggregation Partition the join; plan concurrency; review spill or retry policy
Estimated rows far below actual Poor selectivity estimate Collect column statistics; check predicate form

Where this leads

Reading plans by hand is the diagnostic step. The next step is making the benchmark notice when a plan changes for the worse without anyone looking, which is the subject of Part 4.

Further reading

Performance engineering

Try “evaluation”, “red-teaming”, “governance” or “agents”.