How a SQL Query Actually Executes - Parse, Rewrite, Plan, Execute, Volcano vs Vectorized & Reading Plan Trees
Interview hub: parse -> analyze -> rewrite -> plan -> execute; Volcano iterator vs vectorized vs compiled execution (runnable model counting next() calls, LIMIT pipelining); real PostgreSQL 17 run where one query shape gets index+nested-loop vs seq-scan+hash-join plans depending on the constant; reading plans top-down (control) vs bottom-up (data), inclusive vs exclusive time and q-error (runnable); PG vs MySQL vs SQLite vs DuckDB planner contrasts.
- 1Gist
- 2Maps
- 3Q&A
- 4Sandbox
Voice readout needs Web Speech Synthesis in this browser.
Question ladder
L1
What does the parser know about your tables?
Answer
Nothing. It checks syntax only; names and types are resolved in the analyze step.
L2
Why does a view add no runtime cost by itself?
Answer
The rewriter replaces the view with its definition, so the planner optimizes the combined query over base tables.
L3
What does the planner actually choose?
Answer
Access paths per table, a join algorithm and order for each join, and sort, aggregate and limit strategies, all by estimated cost.
L4
What is the Volcano model?
Answer
Every operator exposes next() returning one row, and parents pull from children, so plans pipeline naturally.
L5
Why do analytical engines vectorize?
Answer
Returning batches of column values amortizes per-row call overhead and lets the CPU use caches and SIMD.
L6
How do you find the expensive node in EXPLAIN ANALYZE?
Answer
Compute exclusive time: inclusive time per loop times loops, minus the children's totals.
L7
Where do you start when a plan is bad?
Answer
At the lowest node where estimated and actual rows diverge, because every node above it was planned on a wrong input size.
Failure modes
Plan flips with the parameter
A rare value gets an index and nested loop, a common value gets seq scans and a hash join, and a cached plan can serve the wrong one.
Per-loop numbers misread
The inner side of a nested loop shows time and rows per loop, so a 0.78 ms probe with loops=50 really costs about 39 ms.
Benchmarks on tiny data
With 100 rows every plan is a seq scan, so a dev benchmark says nothing about the production plan.
Misconceptions
EXPLAIN cost is milliseconds.
Cost is in arbitrary units anchored at one sequential page read. Compare costs only between plans of the same query on the same server.
The top node of a plan is the slow one.
Times are inclusive of children. The hot node is the one with the largest exclusive time.
LIMIT always makes a query cheap.
Above a blocking node such as Sort or Hash, all input is still consumed before the first row comes out.
Interviewer traps
Running EXPLAIN ANALYZE on an UPDATE in production.
It executes the statement. Wrap it in BEGIN and ROLLBACK, or use a replica.
Blaming the hardware for a plan problem.
Faster hardware buys a constant factor; the right plan buys orders of magnitude.
Design scenario
Same prompt for every reader.
Requirements
Find the cause without a deploy, explain the two plans, and propose a fix that keeps small customers fast.
Failure assumptions
- The application uses prepared statements through its driver.
- The largest customer owns a large share of the rows.
- Statistics were last refreshed before a bulk import.
Constraints
- PostgreSQL 17 primary with one read replica.
- No query hints extension installed.
Prompt
A dashboard query over orders is fast for most customers and times out for your largest customer. Explain what the database is doing and how you would investigate.
API
Which parameters reach the query, and does the driver send them as a prepared statement?
Data
Which statistics describe the customer_id distribution, and how stale are they?
Architecture
Where do you capture the slow plan (auto_explain, replica EXPLAIN ANALYZE), and how do you compare it with the fast one?
Overview
A SQL statement is a description of a result, not a program. Between the text you send and the rows you get back, the database runs a small compiler pipeline: parse the text into a tree, analyze and rewrite it (resolve names, expand views, apply rules), plan it (enumerate physically different ways to compute the same answer and pick the cheapest by an estimated cost), and finally execute the chosen plan tree. Almost every "why is this query slow?" conversation is really about the planner: it chose a join algorithm, a join order, an access path and an aggregation strategy from estimates, and one of those estimates was wrong or one of those choices was impossible because of how the query was written.
This cluster teaches that machinery the way a senior engineer needs it in an interview and on call: what the planner is choosing between, how it prices each choice, why a perfectly reasonable plan becomes terrible overnight, and which rewrites give the optimizer room to work. PostgreSQL 17 is the primary example (every plan shown in this cluster was produced by a real local PostgreSQL 17.11 instance unless explicitly marked illustrative), with MySQL, SQLite and DuckDB contrasts where their design differs in an instructive way.
This page: the pipeline, the Volcano iterator vs vectorized execution models, why the same SQL text gets different plans, and how to read a plan tree top-down vs bottom-up. The five sibling pages go deep on join algorithms, the cost-based optimizer and join ordering, cardinality misestimates and plan regressions, sorting/aggregation/spills/parallelism, and query rewrites and SARGability.
Not re-taught here (cross-linked instead): index types, composite/covering/partial indexes, the basics of EXPLAIN and EXPLAIN ANALYZE, statistics and selectivity fundamentals, MVCC, and storage engine internals. Those pages already exist; this cluster sits one level above them, at the optimizer and executor.
The pipeline, step by step
Flow
- 1
Step 1: Parser - SQL text to raw parse tree (syntax only)
- nextStep 2: Analyzer - resolve tables, columns, types into a Query tree
- 2
Step 2: Analyzer - resolve tables, columns, types into a Query tree
- nextStep 3: Rewriter - expand views, apply rules, row-level security
- 3
Step 3: Rewriter - expand views, apply rules, row-level security
- nextStep 4: Planner - generate paths, estimate rows and cost, keep the cheapest
- 4
Step 4: Planner - generate paths, estimate rows and cost, keep the cheapest
- nextStep 5: Executor - run the plan tree, parents pull rows from children
- 5
Step 5: Executor - run the plan tree, parents pull rows from children
- nextStep 6: Rows to client
- 6
Step 6: Rows to client
- 7
pg_statistic, pg_class: row counts, MCVs, histograms
- feeds estimatesStep 4: Planner - generate paths, estimate rows and cost, keep the cheapest
- 8
GUCs: work_mem, random_page_cost, enable_*
- shape cost and choicesStep 4: Planner - generate paths, estimate rows and cost, keep the cheapest
Lesson map
How a SQL Query Actually Executes - Parse, Rewrite, Plan, Execute, Volcano vs Vectorized & Reading Plan Trees
Interview hub: parse -> analyze -> rewrite -> plan -> execute; Volcano iterator vs vectorized vs compiled execution (runnable model counting next() calls, LIMIT pipelining); real PostgreSQL 17 run where one query shape gets index+nested-loop vs seq-scan+hash-join plans depending on the constant; reading plans top-down (control) vs bottom-up (data), inclusive vs exclusive time and q-error (runnable); PG vs MySQL vs SQLite vs DuckDB planner contrasts.
Architecture. Architecture
Select a node to see why it exists, or an edge to see the protocol, direction, effect, and consequence.
Mermaid export
flowchart TB a["Step 1: Parser - SQL text to raw parse tree (syntax only)"] b["Step 2: Analyzer - resolve tables, columns, types into a Query tree"] c["Step 3: Rewriter - expand views, apply rules, row-level security"] d["Step 4: Planner - generate paths, estimate rows and cost, keep the cheapest"] e["Step 5: Executor - run the plan tree, parents pull rows from children"] f["Step 6: Rows to client"] s["pg_statistic, pg_class: row counts, MCVs, histograms"] g["GUCs: work_mem, random_page_cost, enable_*"] a -->|continues| b b -->|continues| c c -->|continues| d d -->|continues| e e -->|continues| f s -->|feeds estimates| d g -->|shape cost and choices| d
- Parse. Pure syntax: keywords, identifiers, literals, operator precedence. The parser does not know whether
ordersexists. A syntax error is the only failure here. - Analyze (semantic analysis). Names are resolved against the catalog, types are assigned,
*is expanded, implicit casts are inserted. This is whereid = '4242'becomesid = 4242(the untyped literal is coerced to the column's type), and whereid::text = '4242'stays a text comparison, which matters a lot later. - Rewrite. Views are replaced by their definitions, rules and row-level security policies are applied. After this step the planner sees base tables only, which is why a view adds no runtime cost by itself and why a predicate on a view can be pushed into the view's tables.
- Plan / optimize. The interesting part. The planner pulls up subqueries and flattens joins where it is semantically safe, then for each base relation generates access paths (Seq Scan, Index Scan, Index Only Scan, Bitmap Heap Scan), for each set of relations generates join paths (Nested Loop, Hash Join, Merge Join, in each order and direction), adds sort/aggregate/limit nodes, and keeps the cheapest path by estimated cost. Estimates come from statistics gathered by
ANALYZEand from cost parameters. - Execute. The chosen plan becomes a tree of executor nodes. In PostgreSQL each node implements an iterator interface; the root asks its child for a row, which asks its child, and so on down to the scans.
- Return. Rows stream to the client as the root produces them. A query with
LIMIT 10on top of a pipelined plan can return after reading a tiny fraction of the table.
Prepared statements split this pipeline: parse/analyze/rewrite happen once at PREPARE, and planning happens either per execution (a custom plan) or once for all parameter values (a generic plan). That split is the root of an entire class of production incidents covered in the misestimates page.
How the optimizer's choices map to the rest of this cluster
| Planner decision | Main inputs | Typical failure | Deep-dive page |
|---|---|---|---|
| Access path per table | Selectivity, index availability, random_page_cost, correlation | Seq scan on a selective predicate because the predicate is not SARGable | Query Rewrites & SARGability |
| Join algorithm per join | Input sizes, equality vs range condition, sort order, work_mem | Nested loop chosen on a "tiny" input that is actually huge | Join Algorithms |
| Join order | Cardinality of every intermediate result | Explodes past 8-12 tables, falls back to GEQO or written order | Cost-Based Optimizer & Join Ordering |
| Row estimates | Statistics, independence assumption, parameter values | Correlated columns, stale stats, generic plans | Cardinality Misestimates & Plan Regressions |
| Sort / aggregate strategy | work_mem, number of groups, LIMIT | Spills to disk, memory blow-ups, missing top-N | Sorting, Aggregation, Spills & Parallel Query |
Row-at-a-time iterators or vectorized batches for a big aggregate?
Prefer
Vectorized batches
next() returns about a thousand values per call, and each operator runs a tight loop.
- The full SUM made 977 calls instead of 1,090,002.
- Same answer, 8550000, from both models.
- Cache and SIMD friendly, which is why DuckDB and ClickHouse use it.
Alternative
Volcano iterators
One row per next() call, pulled through every operator.
- LIMIT 5 stopped after scanning 96 of 1,000,000 rows.
- Simple, composable and quick to the first row, which suits OLTP.
- Per-row calls dominate on large scans.
From SQL text to rows
Diagram 1 condensed. The planner step is where the sibling pages go deep.
- 1
Parse
SQL text becomes a raw parse tree. Only syntax errors fail here. - 2
Analyze and rewrite
Names and types are resolved, views expanded, rules and row-level security applied. - 3
Plan
Paths are generated, rows estimated from statistics, costs compared and the cheapest kept. - 4
Execute
Parents pull rows from children; rows stream to the client. - 5
Estimate goes wrong
A bad row estimate makes the cheapest plan on paper the slowest one in practice.
Two ways to execute a plan: Volcano iterators vs vectorized batches
The classic model, from Goetz Graefe's Volcano work, is the iterator (pull) model: every operator implements open(), next(), close(). A next() call returns one tuple. It is elegant, composable and naturally pipelined: a Limit on top simply stops calling next(), so a filter-and-limit query never touches most of the table.
The cost is per-tuple overhead: a virtual call per row per operator, poor instruction-cache behaviour, and no chance for the CPU to apply SIMD across many values. Analytical engines (MonetDB/X100, Vectorwise, DuckDB, ClickHouse, Snowflake, BigQuery) therefore use vectorized execution: next() returns a batch (DuckDB's default vector is 2,048 values) of one or a few columns, and each operator runs a tight loop over the batch. Another family (HyPer, Umbra, some Spark paths) uses compilation: the plan is turned into machine code with operators fused into loops.
The runnable model below implements a tiny Volcano pipeline with Scan -> Filter -> Project -> Limit, then the same aggregate in batch style, and counts calls.
"""Volcano (iterator) vs vectorized execution, in miniature.
Each operator exposes next(): the parent pulls one row at a time from its child.
We count next() calls to show (a) LIMIT stops the pipeline early and
(b) a vectorized engine makes ~1/batch_size as many calls for the same work.
All table sizes are example values.
"""
calls = {"volcano": 0, "vector": 0}
# ---------- Volcano: one row per next() ----------
class Scan:
def __init__(self, rows): self.rows, self.i = rows, 0
def next(self):
calls["volcano"] += 1
if self.i >= len(self.rows): return None # None = end of stream
r = self.rows[self.i]; self.i += 1; return r
class Filter:
def __init__(self, child, pred): self.child, self.pred = child, pred
def next(self):
calls["volcano"] += 1
while (r := self.child.next()) is not None:
if self.pred(r): return r # pass matching rows upward
return None
class Project:
def __init__(self, child, fn): self.child, self.fn = child, fn
def next(self):
calls["volcano"] += 1
r = self.child.next()
return None if r is None else self.fn(r)
class Limit:
def __init__(self, child, n): self.child, self.n = child, n
def next(self):
calls["volcano"] += 1
if self.n == 0: return None # stop pulling: the scan never finishes
self.n -= 1; return self.child.next()
rows = [{"id": i, "amount": i % 100} for i in range(1_000_000)] # example table
plan = Limit(Project(Filter(Scan(rows), lambda r: r["amount"] > 90), lambda r: r["id"]), 5)
out = []
while (r := plan.next()) is not None: out.append(r)
print("volcano LIMIT 5 ->", out, "| next() calls:", calls["volcano"], "| rows scanned:", plan.child.child.child.i)
# ---------- Vectorized: next() returns a batch (column chunk) ----------
def vector_sum(rows, batch=1024):
"""SUM(amount) WHERE amount > 90, processed 1024 rows per call (like DuckDB/ClickHouse vectors)."""
total = 0
for start in range(0, len(rows), batch):
calls["vector"] += 1
chunk = [r["amount"] for r in rows[start:start + batch]] # column slice
total += sum(a for a in chunk if a > 90) # tight loop over one column
return total
calls["volcano"] = 0
full = Filter(Scan(rows), lambda r: r["amount"] > 90)
vol_total = 0
while (r := full.next()) is not None: vol_total += r["amount"]
vec_total = vector_sum(rows)
print("full SUM: volcano calls =", calls["volcano"], "| vectorized calls =", calls["vector"],
"| same answer:", vol_total == vec_total, vol_total)Output:
volcano LIMIT 5 -> [91, 92, 93, 94, 95] | next() calls: 112 | rows scanned: 96
full SUM: volcano calls = 1090002 | vectorized calls = 977 | same answer: True 8550000Two lessons hide in that output. First, pipelining: the LIMIT 5 pipeline scanned 96 rows of a million, because nothing above the scan needed more. That is why ORDER BY x LIMIT 10 with an index on x is fast and the same query without the index is not (a sort is a blocking operator: it must consume all input before emitting its first row). Second, per-call overhead: the vectorized loop made about a thousandth as many operator calls for identical work. In a real engine that turns into fewer branch mispredictions, better cache locality on columnar data and SIMD.
| Execution model | How rows move | Strengths | Weaknesses | Where you see it |
|---|---|---|---|---|
| Volcano / iterator (tuple at a time) | Parent calls next(), gets one row | Simple, composable, pipelines naturally, low latency to first row | Per-row virtual calls, poor CPU efficiency for big scans | PostgreSQL, MySQL, SQLite (VDBE bytecode, still row at a time), SQL Server row mode |
| Vectorized (batch at a time) | Parent gets a vector of ~1-2k values per column | Amortizes call overhead, cache and SIMD friendly, great for scans and aggregates | More complex operators, more memory per operator, less natural for OLTP point lookups | DuckDB, ClickHouse, Vectorwise, SQL Server batch mode, Databricks Photon |
| Compiled (data-centric codegen) | Operators fused into generated loops | Near hand-written speed, values stay in registers | Compilation latency, harder to debug, complex engine | HyPer, Umbra, Spark whole-stage codegen; PostgreSQL JIT compiles expressions only |
What happens if you choose otherwise. If you push a large analytical workload (scan 500M rows, group by a few keys) through a row-at-a-time OLTP engine, you pay per-row interpretation costs that a vectorized columnar engine avoids, which is why teams offload dashboards to DuckDB, ClickHouse or a warehouse instead of adding indexes forever. If you put point-lookup OLTP traffic on a columnar engine, every single-row update touches many column files and batch machinery that has nothing to amortize. PostgreSQL narrows the gap with parallel query (see the sorting/parallel page) and JIT-compiled expressions, but its executor remains an iterator engine.
Why the same query gets different plans
The planner does not choose a plan for the SQL text. It chooses for the text plus the parameter values plus the current statistics plus the cost settings plus the available indexes. Change any of them and the cheapest plan can change. That is a feature, not a bug, and it is also why plans "regress" without a deploy.
The real PostgreSQL session below runs one query shape twice with two different constants. status = 'refunded' matches 50 of 200,000 orders; status = 'shipped' matches about 90%.
-- Page 1: the same query shape gets different plans when the predicate's selectivity changes.
-- Self-contained: builds its own schema p1 with deterministic data (no random()).
\pset footer off
SET client_min_messages = warning;
DROP SCHEMA IF EXISTS p1 CASCADE;
CREATE SCHEMA p1;
SET search_path = p1;
SET max_parallel_workers_per_gather = 0; -- keep plans small and readable for the lesson
-- 20k customers, 200k orders; status is skewed: 'shipped' is common, 'refunded' is rare
CREATE TABLE customers AS
SELECT g AS id, (ARRAY['US','IN','DE','BR','JP'])[1 + (g / 7) % 5] AS country
FROM generate_series(1, 20000) g;
ALTER TABLE customers ADD PRIMARY KEY (id);
CREATE TABLE orders AS
SELECT g AS id,
1 + (g * 7919) % 20000 AS customer_id,
CASE WHEN g % 4000 = 0 THEN 'refunded' WHEN g % 10 = 0 THEN 'pending' ELSE 'shipped' END AS status,
(g % 500) + 0.99 AS amount
FROM generate_series(1, 200000) g;
ALTER TABLE orders ADD PRIMARY KEY (id);
CREATE INDEX orders_status_idx ON orders (status);
ANALYZE customers; ANALYZE orders;
-- Same SQL text, different constant: rare value -> index path + nested loop; common value -> seq scans + hash join
EXPLAIN (COSTS OFF)
SELECT c.country, count(*) FROM orders o JOIN customers c ON c.id = o.customer_id
WHERE o.status = 'refunded' GROUP BY c.country;
EXPLAIN (COSTS OFF)
SELECT c.country, count(*) FROM orders o JOIN customers c ON c.id = o.customer_id
WHERE o.status = 'shipped' GROUP BY c.country;
-- Read a full plan with estimates vs actuals (timing off so the output is stable)
EXPLAIN (ANALYZE, TIMING OFF, SUMMARY OFF)
SELECT c.country, count(*) FROM orders o JOIN customers c ON c.id = o.customer_id
WHERE o.status = 'refunded' GROUP BY c.country;Real output (PostgreSQL 17.11, local sandbox, psql; setup DDL echoes omitted):
SET max_parallel_workers_per_gather = 0;
EXPLAIN (COSTS OFF)
SELECT c.country, count(*) FROM orders o JOIN customers c ON c.id = o.customer_id
WHERE o.status = 'refunded' GROUP BY c.country;
QUERY PLAN
------------------------------------------------------------------
GroupAggregate
Group Key: c.country
-> Sort
Sort Key: c.country
-> Nested Loop
-> Index Scan using orders_status_idx on orders o
Index Cond: (status = 'refunded'::text)
-> Index Scan using customers_pkey on customers c
Index Cond: (id = o.customer_id)
EXPLAIN (COSTS OFF)
SELECT c.country, count(*) FROM orders o JOIN customers c ON c.id = o.customer_id
WHERE o.status = 'shipped' GROUP BY c.country;
QUERY PLAN
--------------------------------------------------
HashAggregate
Group Key: c.country
-> Hash Join
Hash Cond: (o.customer_id = c.id)
-> Seq Scan on orders o
Filter: (status = 'shipped'::text)
-> Hash
-> Seq Scan on customers c
EXPLAIN (ANALYZE, TIMING OFF, SUMMARY OFF)
SELECT c.country, count(*) FROM orders o JOIN customers c ON c.id = o.customer_id
WHERE o.status = 'refunded' GROUP BY c.country;
QUERY PLAN
-------------------------------------------------------------------------------------------------------------------------------
GroupAggregate (cost=286.69..286.98 rows=5 width=11) (actual rows=4 loops=1)
Group Key: c.country
-> Sort (cost=286.69..286.77 rows=33 width=3) (actual rows=50 loops=1)
Sort Key: c.country
Sort Method: quicksort Memory: 25kB
-> Nested Loop (cost=0.58..285.85 rows=33 width=3) (actual rows=50 loops=1)
-> Index Scan using orders_status_idx on orders o (cost=0.29..51.79 rows=33 width=4) (actual rows=50 loops=1)
Index Cond: (status = 'refunded'::text)
-> Index Scan using customers_pkey on customers c (cost=0.29..7.09 rows=1 width=7) (actual rows=1 loops=50)
Index Cond: (id = o.customer_id)Read the difference:
- Rare value: the planner expects about 33 rows (the truth is 50), so it uses the
statusindex, then for each of those rows probescustomers_pkey(Nested Loop with an inner Index Scan,loops=50). Sorting 50 rows for a GroupAggregate is trivial. - Common value: about 180,000 rows survive the filter, so random index probes would be far more expensive than reading both tables sequentially. The planner switches to Seq Scans, a Hash Join (build a hash table on the 20,000 customers, stream orders through it) and a HashAggregate (no sort).
- The EXPLAIN ANALYZE output shows estimated vs actual rows on every node (
rows=33vsactual rows=50). When those two numbers are close, the planner had a fair chance. When they differ by orders of magnitude, start there (the misestimates page is about exactly that).
| What changed | Example | Typical plan effect |
|---|---|---|
| Parameter value | status = 'refunded' vs 'shipped' | Index + nested loop vs seq scan + hash join |
| Statistics | Table grew 100x since last ANALYZE, new values outside the histogram | Estimates too low, nested loops on big inputs |
| Cost settings | random_page_cost 4 -> 1.1, work_mem 4MB -> 64MB | Index scans become attractive; hash/sort stop spilling or get chosen |
| Schema | New index, dropped index, new extended statistics | New access paths, better or worse estimates |
| Prepared statement | Custom plan vs generic plan after 5 executions | One plan for every value, chosen without seeing the value |
| Engine version | Planner improvements (Memoize in PG 14, incremental sort in PG 13) | New node types appear, occasionally a regression |
Reading a plan tree: top-down for control, bottom-up for data
EXPLAIN prints the root first, and each -> is a child. Two reading orders are both correct, for different questions:
- Top-down (control flow): "what does the client ultimately wait on?" The root pulls from its children. A
Limitat the top can stop everything below it early; aSortorHashbelow it must finish completely before the parent sees a row. - Bottom-up (data flow): "where do rows come from and how do they grow or shrink?" Start at the most indented leaves (scans), follow rows upward through joins, aggregates and sorts. Check estimated vs actual rows at each step: the first node (from the bottom) where they diverge is usually the root cause, because every node above it was planned on a wrong input size.
Two details trip up even experienced engineers. EXPLAIN ANALYZE actual time is inclusive of children and reported per loop, so the inner side of a nested loop that shows actual time=0.78 with loops=50 really costs about 39 ms in total. And the "hottest" node is the one with the largest exclusive time, not the one at the top.
// Reading a plan tree: control flows top-down (parents pull), data flows bottom-up (children produce).
// EXPLAIN ANALYZE times are INCLUSIVE of children and PER LOOP, so "where did the time go?"
// needs exclusive time = (node time - children time) * loops. Times below are example values.
interface PlanNode {
op: string;
totalMsPerLoop: number; // like "actual time=..X" (end time per loop), inclusive of children
loops: number;
estRows: number;
actualRowsPerLoop: number;
children: PlanNode[];
}
const plan: PlanNode = {
op: "GroupAggregate", totalMsPerLoop: 41.0, loops: 1, estRows: 5, actualRowsPerLoop: 4,
children: [{
op: "Sort", totalMsPerLoop: 40.6, loops: 1, estRows: 53, actualRowsPerLoop: 50,
children: [{
op: "Nested Loop", totalMsPerLoop: 40.2, loops: 1, estRows: 53, actualRowsPerLoop: 50,
children: [
{ op: "Index Scan orders_status_idx", totalMsPerLoop: 0.4, loops: 1, estRows: 53, actualRowsPerLoop: 50, children: [] },
{ op: "Index Scan customers_pkey", totalMsPerLoop: 0.78, loops: 50, estRows: 1, actualRowsPerLoop: 1, children: [] },
],
}],
}],
};
// Top-down print, like EXPLAIN output (indentation = depth).
function print(n: PlanNode, depth = 0): void {
const pad = depth === 0 ? "" : " ".repeat(depth) + "-> ";
console.log(`${pad}${n.op} (est rows=${n.estRows}, actual rows=${n.actualRowsPerLoop} loops=${n.loops})`);
n.children.forEach((c) => print(c, depth + 1));
}
// Post-order walk = the order in which nodes *produce their first rows* (leaves start first).
function postOrder(n: PlanNode, out: string[] = []): string[] {
n.children.forEach((c) => postOrder(c, out));
out.push(n.op);
return out;
}
// Exclusive time: subtract children's inclusive totals (each child's per-loop time * its loops).
function exclusive(n: PlanNode, acc: Array<[string, number]> = []): Array<[string, number]> {
const childTotal = n.children.reduce((s, c) => s + c.totalMsPerLoop * c.loops, 0);
acc.push([n.op, +(n.totalMsPerLoop * n.loops - childTotal).toFixed(2)]);
n.children.forEach((c) => exclusive(c, acc));
return acc;
}
// Misestimate ratio (q-error) per node: max(est/actual, actual/est) using total actual rows.
function qErrors(n: PlanNode, acc: Array<[string, number]> = []): Array<[string, number]> {
const actual = Math.max(1, n.actualRowsPerLoop * n.loops);
const est = Math.max(1, n.estRows * n.loops);
acc.push([n.op, +Math.max(est / actual, actual / est).toFixed(2)]);
n.children.forEach((c) => qErrors(c, acc));
return acc;
}
print(plan);
console.log("production order (bottom-up):", postOrder(plan).join(" => "));
const ex = exclusive(plan).sort((a, b) => b[1] - a[1]);
console.log("exclusive ms, hottest first:", ex.map(([o, t]) => `${o}=${t}`).join(", "));
console.log("q-error per node:", qErrors(plan).map(([o, q]) => `${o}=${q}`).join(", "));Output:
GroupAggregate (est rows=5, actual rows=4 loops=1)
-> Sort (est rows=53, actual rows=50 loops=1)
-> Nested Loop (est rows=53, actual rows=50 loops=1)
-> Index Scan orders_status_idx (est rows=53, actual rows=50 loops=1)
-> Index Scan customers_pkey (est rows=1, actual rows=1 loops=50)
production order (bottom-up): Index Scan orders_status_idx => Index Scan customers_pkey => Nested Loop => Sort => GroupAggregate
exclusive ms, hottest first: Index Scan customers_pkey=39, Nested Loop=0.8, GroupAggregate=0.4, Sort=0.4, Index Scan orders_status_idx=0.4
q-error per node: GroupAggregate=1.25, Sort=1.06, Nested Loop=1.06, Index Scan orders_status_idx=1.06, Index Scan customers_pkey=1ExpectedGroupAggregate (est rows=5, actual rows=4 loops=1) -> Sort (est rows=53, actual rows=50 loops=1) -> Nested Loop (est rows=53, actual rows=50 loops=1) -> Index Scan orders_status_idx (est rows=53, actual rows=50 loops=1) -> Index Scan customers_pkey (est rows=1, actual rows=1 loops=50) production order (bottom-up): Index Scan orders_status_idx => Index Scan customers_pkey => Nested Loop => Sort => GroupAggregate exclusive ms, hottest first: Index Scan customers_pkey=39, Nested Loop=0.8, GroupAggregate=0.4, Sort=0.4, Index Scan orders_status_idx=0.4 q-error per node: GroupAggregate=1.25, Sort=1.06, Nested Loop=1.06, Index Scan orders_status_idx=1.06, Index Scan customers_pkey=1
Press Run. Snippets must be self-contained — no network, files, or native modules.
The times in that snippet are example values, but the arithmetic is exactly what tools like explain.depesz.com and pgMustard do for you: subtract children, multiply by loops, rank. The q-error column (max of est/actual and actual/est) is the standard measure from the optimizer literature for "how wrong was the estimate".
Flow
- 1
Step 4: GroupAggregate - root, client waits here
- nextStep 3: Sort - blocking, must see all input first
- 2
Step 3: Sort - blocking, must see all input first
- nextStep 2: Nested Loop - for each outer row, run inner
- 3
Step 2: Nested Loop - for each outer row, run inner
- nextStep 1a: Index Scan orders_status_idx - outer, runs once
- nextStep 1b: Index Scan customers_pkey - inner, runs loops=50
- 4
Step 1a: Index Scan orders_status_idx - outer, runs once
- rows flow upStep 2: Nested Loop - for each outer row, run inner
- 5
Step 1b: Index Scan customers_pkey - inner, runs loops=50
- rows flow upStep 2: Nested Loop - for each outer row, run inner
Engine contrasts worth knowing
- PostgreSQL: exhaustive dynamic programming over join orders up to
geqo_threshold(default 12) FROM items, genetic optimizer beyond; three join algorithms (nested loop, hash, merge) plus Memoize; no built-in hints (pg_hint_plan is an extension); plan caching only for prepared statements. - MySQL 8.x: historically nested-loop centric (with Block Nested Loop), hash join since 8.0.18 (replacing BNL), no merge join; optimizer hints in comments (
/*+ JOIN_ORDER(...) */);EXPLAIN FORMAT=TREEandEXPLAIN ANALYZEgive a Volcano-style tree. - SQLite: the "Next Generation Query Planner" (3.8+) searches join orders with a bounded best-N path search; nested loops only (with automatic transient indexes);
EXPLAIN QUERY PLANgives a compact access-path view (used in the rewrites page). - DuckDB: vectorized, columnar, parallel by default, hash joins and hash aggregates everywhere, a join-order optimizer based on dynamic programming over the join graph;
EXPLAIN ANALYZEprints per-operator timings.
What happens if you choose otherwise
- Treat the plan as fixed: you will be surprised when a stats refresh, a data-skew change or a version upgrade flips the plan. Treat plans as a function of data, and monitor for plan changes on critical queries (auto_explain, pg_stat_statements).
- Read only the top line of EXPLAIN: you see total cost, but not that one inner index scan ran 2 million loops. Always read bottom-up for row counts.
- Benchmark on an empty dev database: with 100 rows every plan is a seq scan; production with 100 million rows uses a completely different plan. Test with production-like volume and distribution, or at least copy the statistics.
- Force plans with enable_ switches in production:* they are blunt instruments meant for diagnosis (as used in this cluster's demos), not configuration. They affect every query in the session.
Pitfalls
- Confusing cost units with milliseconds. Cost is in arbitrary units anchored at
seq_page_cost = 1.0. Compare costs only between plans of the same query on the same server. - Running
EXPLAIN ANALYZEon aDELETE/UPDATEin production: it executes the statement. Wrap inBEGIN; ... ROLLBACK;. - Forgetting that
EXPLAINwithoutANALYZEshows estimates only. Two plans with identical estimates can have wildly different real behaviour. - Ignoring loops: per-loop numbers on the inner side of nested loops hide the real total.
- Blaming the executor for a plan problem: if the plan is wrong, faster hardware buys you a constant factor; the right plan buys orders of magnitude.
How the code was checked
- The Volcano model ran under Python 3.13 and the plan-tree reader under
tsc --strictand Node 22. Their output blocks are the real captured output. - The SQL script ran through psql against a local PostgreSQL 17.11 sandbox. The plans are real output; only the setup DDL echoes were omitted.
- Table sizes in the Python model and the per-node times in the TypeScript model are example values, labeled in the code comments.
The cluster map
- Join Algorithms - Nested Loop vs Index Nested Loop vs Hash Join vs Merge Join: nested loop, index nested loop, hash join and merge join, their cost formulas, real PostgreSQL plans for each (including Memoize and a hash join that spills to 8 batches), and the planned-10-got-1M failure.
- Cost-Based Optimizer & Join Ordering - Cost Constants, Dynamic Programming, Left-Deep vs Bushy, join_collapse_limit & GEQO: PostgreSQL cost constants checked against pg_class, random_page_cost flipping seq vs index scan, Selinger dynamic programming, left-deep vs bushy trees, join_collapse_limit and GEQO.
- Cardinality Misestimates & Plan Regressions - Correlated Columns, Extended Statistics, Generic Plans & the Slow-Overnight Runbook: the independence assumption, extended statistics, stale stats, custom vs generic plans for skewed prepared statements, auto_explain and pg_stat_statements, and the slow-overnight runbook.
- Sorting, Aggregation, Spills & Parallel Query - External Merge, Top-N, HashAggregate vs GroupAggregate & work_mem Math: quicksort, external merge, top-N heapsort and incremental sort, HashAggregate vs GroupAggregate, real spills, a Gather plan, and work_mem worst-case sizing.
- Query Rewrites & SARGability - Unnesting, NOT IN vs NOT EXISTS, OR to UNION, Keyset Pagination & ORM N+1: subquery unnesting, the NOT IN NULL trap, SARGable predicates, BitmapOr vs UNION, OFFSET vs keyset pagination, and ORM N+1 detection.
Interview Q&A
Walk me through what happens when PostgreSQL receives a SELECT.
Answer
Parse to a raw tree (syntax only), analyze to a Query tree with resolved names and types, rewrite (views, rules, RLS), plan (generate access and join paths, estimate rows from statistics, cost them, keep the cheapest), then execute the plan tree with iterator nodes that pull rows from their children, streaming results to the client.
Why might the same query be fast for one customer and slow for another?
Answer
The plan depends on the parameter's selectivity. A rare value favours an index and nested loop; a common one favours seq scans and hash joins. With prepared statements and a generic plan, both customers get the same plan, which can be wrong for the outlier. Check estimated vs actual rows and whether a generic plan is in use.
What is the Volcano model and why do analytical engines move away from it?
Answer
Every operator exposes next() returning one tuple; parents pull from children. It pipelines well and is simple, but per-tuple virtual calls and interpretation overhead dominate on large scans. Vectorized engines return batches of column values to amortize overhead and use SIMD; compiled engines fuse operators into generated loops.
How do you find the expensive node in an EXPLAIN ANALYZE plan?
Answer
Compute exclusive time per node: inclusive time per loop times loops, minus children's totals. Then look for the lowest node where estimated and actual rows diverge, because everything above it was planned on bad input.
Which operators are blocking, and why does that matter for LIMIT?
Answer
Sort, Hash (build side), HashAggregate and Materialize must consume all their input before producing output (a top-N sort keeps only N but still reads everything). A LIMIT above a fully pipelined path can stop early; above a blocking node it saves only the output, not the input work.
What does the rewriter do with a view?
Answer
It replaces the view reference with the view's query, so the planner optimizes the combined query; predicates can be pushed into the view's tables. Views add no execution cost by themselves (materialized views are different: they are real tables).
What changes the plan for the same SQL text without a deploy?
Answer
Parameter values, refreshed or stale statistics, cost settings such as random_page_cost and work_mem, new or dropped indexes, the switch to a generic plan for a prepared statement, and engine upgrades.
How does a prepared statement split the pipeline?
Answer
Parse, analyze and rewrite happen once at PREPARE. Planning happens per execution for a custom plan, or once for all values for a generic plan, which is the root of a whole class of skew incidents.
How do PostgreSQL and DuckDB differ in execution?
Answer
PostgreSQL runs a row-at-a-time iterator engine with optional parallel workers and JIT-compiled expressions. DuckDB is columnar, vectorized and parallel by default, so it wins large scans and aggregates and is not built for many small concurrent writes.
Check yourself
Take one slow query from your own system and run EXPLAIN (ANALYZE, BUFFERS) on a replica. Mark the lowest node where estimated and actual rows differ by more than 10x, then compute the exclusive time of the three hottest nodes by hand.
Elsewhere in the library
These pages stay as they are. This lesson only points at them: Indexes, Cardinality & EXPLAIN Plans, EXPLAIN & EXPLAIN ANALYZE, Database Storage Engines — WAL, B-Trees & LSM Trees, SQL Analytics — Window Functions, CTEs & Set-Based Thinking, Performance Engineering — Profiling, Flame Graphs, Latency Budgets & Hot Paths.
Go Deeper
- PostgreSQL docs: The Path of a Query (parser, rewriter, planner, executor)
- PostgreSQL docs: Planner/Optimizer overview
- PostgreSQL docs: Using EXPLAIN
- Boncz, Zukowski, Nes: MonetDB/X100 - Hyper-Pipelining Query Execution (CIDR 2005), the vectorized execution paper
- Leis et al.: How Good Are Query Optimizers, Really? (VLDB 2015)
- CMU 15-445/645 Database Systems schedule (query execution and optimization lectures)
- CMU Database Group on YouTube (Andy Pavlo's lectures)
- DuckDB: Why DuckDB (vectorized, columnar, in-process)
- SQLite: The Next-Generation Query Planner
- explain.depesz.com: paste a plan, get exclusive times per node