skip to content

How do you read an execution plan tree - which operator runs first and how does data flow between nodes - and how do you interpret a node reported as '8 rows, 40,000 executions'?

level: middleimportance: must knowfreq 68%

answer

  1. leaves first, rows flow upward
  2. nested loop: outer once, inner per outer row
  3. rows x executions = real total
  4. blocking operators delay first row
  5. parent time is cumulative - subtract children

basics

~20 s

Read innermost/deepest nodes first: leaves produce rows, parents consume them, results flow upward to the root. A node showing 8 rows over 40,000 executions was started 40,000 times and produced about 8 rows each time - roughly 320,000 rows in total, so multiply before judging cost.

solid answer

~60 s

A plan is a tree printed with the root - the operator that returns the final result - at the top and its inputs indented beneath it. Execution begins at the **leaves** (the access methods), and rows flow **upward**: each operator pulls rows from its children, transforms them and emits rows to its parent. So you read the *structure* top-down to understand intent, but you read the *work* bottom-up. Order among siblings matters. For a nested loop, the first child is the outer/driving input and the second child is executed once per outer row. For a hash join, the build input is consumed completely before the probe input streams through. Many engines report the inner side of a loop as a **per-execution average plus a loop count**. '8 rows, 40,000 executions' means the operator ran 40,000 times and averaged 8 rows, so about 320,000 rows and 40,000 index descents in total. Timings are often per-execution too, so a node showing 0.4 ms is really 16 seconds of work. Averages also hide skew: one execution may have done most of it.

code

text · 6 lines
text
Nested Loop  (est rows=120)  actual rows=310,522  time=18,402ms
  -> Table Scan on orders  (est rows=40,000)  actual rows=40,118  loops=1  time=310ms
  -> Index Range Scan on order_items (order_id = o.id)
         (est rows=3)  actual rows=8  loops=40,118  time=0.44ms

# inner side total = 8 x 40,118 ~ 320k rows, 0.44ms x 40,118 ~ 17.6s

go deeper

for a junior

Say that leaves run first and rows flow upward to the root, and that a node's row count is what it emits to its parent.

for a middle

Add the loop-count multiplication, the outer/inner asymmetry of nested loops, and the fact that parent times are usually cumulative.

for a senior

Turn it into a routine: rows touched vs returned, multiply loops out, compute exclusive time, then find the lowest node where estimate and actual diverge.

for a principal

Discuss blocking vs pipelined execution and its effect on first-row latency and streaming workloads, plus the conventions that make parallel plans easy to misread.

## The shape of the tree Plans are printed as an indented tree. The least-indented line is the **root**: the operator that hands the final rows to the client. Everything indented under a node is its **input**. Leaves are access methods - table scans, index scans, lookups - and they are where rows enter the plan. Data flows from leaves toward the root. That is why the reading rule is 'deepest first': the innermost operator produces rows, its parent consumes and transforms them, and so on up. The root is where you find out what the query returns; the leaves are where you find out what it touched. ``` Limit -> Sort (order by o.created_at) -> Hash Join (o.customer_id = c.id) -> Table Scan on orders -> Hash Build -> Index Range Scan on customers ``` The first physical work here is the index range scan on customers plus the table scan on orders; the last is the limit. ## Sibling order is not decorative - **Nested loop**: the first (outer) input runs once and drives the join; the second (inner) input is re-executed for each outer row. If the outer side yields 40,000 rows, the inner subtree runs 40,000 times. - **Hash join**: the build side is consumed fully into a hash table first, then the probe side streams past it. Build-side size is what determines memory. - **Merge join**: both inputs must arrive sorted on the join key; a sort may be injected under either side. A related distinction is **pipelined vs blocking** operators. A filter or a nested loop is pipelined - it can emit its first row almost immediately, which is what makes a top-N query with a matching index fast. A sort, a hash build, or a full aggregate is **blocking**: it must read all of its input before emitting anything. When a query is slow to return its first row, look for the lowest blocking operator. ## Reading the per-node numbers Each node reports some subset of: estimated rows, actual rows, executions (loops), time, memory. The single most common misreading is treating a per-execution average as a total. If a node shows **rows=8, executions=40,000**, the operator was started 40,000 times and averaged 8 rows per start. Total rows are roughly 320,000, and if it was an index lookup, that is 40,000 separate index descents. Engines differ in whether time is also averaged; when it is, a node reporting 0.4 ms per execution over 40,000 executions accounts for about 16 seconds. Always establish which convention the output uses before drawing conclusions, then multiply. Averages also conceal **skew**. Forty thousand executions averaging 8 rows might be 39,999 executions returning 1 row and a single execution returning 300,000. Where an engine reports minimum/maximum rows or per-worker figures, check them; where it does not, be suspicious of a large loop count over a column you know is unevenly distributed. ## Attributing time Most plan formats report **cumulative** time at each node: a parent's elapsed time includes all of its children's. To find where time is actually spent, subtract the children's totals from the parent's - the operator with the largest *exclusive* time is the one to attack. A root node showing 12 seconds tells you nothing on its own; a hash aggregate showing 11.5 seconds cumulative over children totalling 0.3 seconds tells you everything. Parallel plans add another wrinkle: per-worker rows may be reported per worker or summed, and wall-clock time is not the sum of worker times. Read the convention before concluding. ## A practical routine 1. Find the root and note the final row count. 2. Walk to the leaves and note what each access method touched. 3. Compute rows touched versus rows returned. A large ratio points at the access path or a filter applied too late. 4. Multiply out any node with a loop count. 5. Compute exclusive time per node and rank. 6. Only then form a hypothesis, and check it against the estimate-versus-actual gap at the lowest node where the two diverge.

  • A plan's root node reports 12 seconds. How do you find the operator actually responsible?
    Most plan formats report cumulative time, so the root's 12 seconds includes every descendant. Subtract each node's children's totals from its own to get exclusive time, and rank by that. Remember to multiply per-execution timings by their loop count first, otherwise a repeatedly executed inner node looks trivial.
  • Why does a query with a LIMIT sometimes return instantly and sometimes take as long as the unlimited query?
    It depends on whether the operators beneath the limit are pipelined or blocking. If an index already supplies the required order, rows stream up and the limit stops early. If a sort, hash build or full aggregate sits underneath, that operator must consume its entire input before emitting anything, so the limit saves almost nothing.

saying these in an interview costs you the question

  • Reading a plan strictly top-down as the execution order
  • Treating a per-execution row or time figure as the operator's total
  • Ignoring which side of a nested loop is the driving input
  • Adding up cumulative node times and concluding the query took far longer than it did
  • Assuming a LIMIT always makes a query cheap regardless of what is beneath it

context