SQL is declarative: you say what rows you want, not how to find them. The query processor decides the how. For a single query with a few joins there can be thousands of possible execution strategies, and the fastest can be a million times quicker than the slowest. This lesson follows a query through the pipeline (parse, rewrite, optimise, execute), shows how relational algebra trees represent plans, compares the algorithms for selections and joins with worked I/O cost calculations, explains heuristic and cost-based optimisation, and teaches you to read EXPLAIN output. It finishes with the slow-query fixes you will actually use at work. Interviewers ask about join algorithms and their costs, about why the optimiser ignored an index, and how you would debug a slow query.
You should already know indexing and B+ trees and basic SQL.
The steps of query processing
SQL text
|
v
+---------+ parse tree +-----------+ logical plan +-----------+
| Parser | -----------> | Rewriter | -------------> | Optimiser |
+---------+ +-----------+ +-----------+
syntax, names, views expanded, |
types checked simplified physical plan |
v
result rows <-- +-----------+
| Executor |
+-----------+
- Parsing. The parser checks the SQL grammar and builds a parse tree. The binder (or analyser) then resolves names against the catalog (the database's metadata about tables, columns and types): does table
ordersexist, istotala column, is comparing it to'abc'a type error? Errors such as "column does not exist" come from here. - Rewriting. The query is transformed into an equivalent, simpler logical form: views are replaced by their definitions, row-level security predicates are added, some subqueries are turned into joins (for example
IN (SELECT ...)into a semi-join), constants are folded (1 + 1becomes2), and redundant conditions are removed. - Optimisation. The optimiser enumerates alternative physical plans (which access method for each table, which join algorithm, which join order) and picks the one with the lowest estimated cost.
- Execution. The executor runs the plan. Most databases use the iterator model (also called the Volcano model): each operator has
open(),next()andclose()methods, and each call tonext()pulls one row from its child. Rows flow up the tree without being written to disk in between, which is called pipelining. Some operators, such as sort and the build side of a hash join, must consume their whole input first; they are blocking operators.
Many databases also cache plans for prepared statements so steps 1 to 3 are skipped on repeated execution.
Relational algebra trees
The optimiser works on a tree of relational algebra operators. The basic ones:
| Symbol | Name | SQL equivalent |
|---|---|---|
| σ (sigma) | Selection: keep rows matching a condition | WHERE |
| π (pi) | Projection: keep some columns | the SELECT list |
| ⋈ (bowtie) | Join | JOIN ... ON |
| × | Cartesian product | CROSS JOIN |
| γ (gamma) | Grouping and aggregation | GROUP BY |
Take this query on an orders schema:
SELECT c.name, o.total
FROM customers c
JOIN orders o ON o.customer_id = c.id
WHERE c.city = 'Pune' AND o.status = 'pending';
A direct translation (join everything, then filter, then project):
π name, total
|
σ city='Pune' AND status='pending'
|
⋈ o.customer_id = c.id
/ \
customers orders
An equivalent, much cheaper tree pushes the filters down to the tables, so the join processes far fewer rows:
π name, total
|
⋈ o.customer_id = c.id
/ \
σ city='Pune' σ status='pending'
| |
customers orders
Both trees are logically the same: they produce the same result. A physical plan goes further and annotates each node with an algorithm: "index scan on orders using idx_status", "hash join with customers as the build side".
Cost model basics
To compare plans the optimiser needs a cost model. In textbooks, the cost is the number of disk block transfers (and sometimes seeks), because disk I/O traditionally dominated. Real systems add CPU cost per row and per operator evaluation, and weight random I/O higher than sequential I/O.
Notation used in this lesson:
| Symbol | Meaning |
|---|---|
| n_r | number of tuples (rows) in relation r |
| b_r | number of blocks (pages) holding r |
| M | number of memory blocks (buffer pages) available to the operator |
| h | height of a B+ tree index (levels from root to leaf) |
Costs ignore writing the final result, since every plan must produce the same output.
Selection algorithms
How to evaluate σ condition (r):
| Algorithm | When it applies | Cost (block transfers) |
|---|---|---|
| Linear scan | Always | b_r (or b_r / 2 on average for equality on a unique key, since you can stop at the match) |
| Binary search | File sorted on the attribute, equality | about ⌈log₂ b_r⌉ plus the blocks with matching rows |
| Primary (clustered) index, equality on key | B+ tree on a unique key | h + 1 |
| Primary index, equality on non-key | Clustered on the attribute | h + number of blocks with matches (they are contiguous) |
| Secondary index, equality | Non-clustered index | h + one block per matching row (each match may be on a different block) |
Primary index, range (>, BETWEEN) | Clustered | h + blocks in the range |
| Secondary index, range | Non-clustered | h + leaf blocks scanned + one block per matching row |
Worked example. orders has 10,000 rows in 400 blocks (25 rows per block). A secondary B+ tree index on status has height 3.
- Query
WHERE status = 'pending'matches 50 rows. Index: 3 + 50 = 53 block reads (worst case, each row in a different block). Linear scan: 400. Use the index. - Query
WHERE status = 'shipped'matches 9,000 rows. Index: 3 + 9,000 = 9,003 random reads. Linear scan: 400 sequential reads. Use the scan.
This is the single most important intuition about indexes: a secondary index wins only when the condition is selective (matches a small fraction of rows). Past a few per cent of the table, a sequential scan is usually cheaper. This is why the optimiser sometimes "ignores your index", and it is usually right.
Complex conditions. For A AND B, use the index for the most selective condition and check the rest on the fetched rows, or intersect row-id lists from two indexes (PostgreSQL's bitmap AND). For A OR B, use the union of row-id lists from two indexes, or fall back to a scan if either condition has no index.
Join algorithms with worked costs
Joins are where plans differ most. We use one running example throughout and compute every number.
orders(call it r): n_r = 10,000 rows, b_r = 400 blocks.customers(call it s): n_s = 5,000 rows, b_s = 100 blocks.- The join condition is
orders.customer_id = customers.id.
1. Nested loop join
For every row of the outer relation, scan the whole inner relation.
for each row tr in r: -- outer
for each row ts in s: -- inner
if tr.customer_id = ts.id: output (tr, ts)
Worst case (only one block of each relation in memory): the inner relation is read once per outer row.
- Cost = n_r × b_s + b_r.
- orders outer: 10,000 × 100 + 400 = 1,000,400 block transfers.
- customers outer: 5,000 × 400 + 100 = 2,000,100 block transfers.
If the smaller relation fits entirely in memory, read each relation once: 400 + 100 = 500.
Nested loop works for any join condition (including <, LIKE, or arbitrary functions), which is why databases keep it, but without an index or memory it is the slowest option.
2. Block nested loop join
Improve by looping over blocks instead of rows: for each block of the outer, scan the inner once and compare every pair of rows in memory. With M buffer blocks, use M − 2 blocks for a chunk of the outer, one for the inner, one for output.
- Worst case (M = 3, one outer block at a time): cost = b_r × b_s + b_r.
- orders outer: 400 × 100 + 400 = 40,400.
- customers outer: 100 × 400 + 100 = 40,100.
- General: cost = ⌈b_outer / (M − 2)⌉ × b_inner + b_outer.
- With M = 12 and customers outer: ⌈100 / 10⌉ × 400 + 100 = 10 × 400 + 100 = 4,100.
- With M = 12 and orders outer: ⌈400 / 10⌉ × 100 + 400 = 40 × 100 + 400 = 4,400.
Rule: put the smaller relation on the outside, and give the operator as much memory as possible.
3. Index nested loop join
If the inner relation has an index on the join column, replace the inner scan with an index lookup per outer row.
- Cost = b_r + n_r × c, where c is the cost of one index lookup plus fetching the matching row(s).
- customers has a B+ tree on its primary key
idwith height 3, so c = 3 + 1 = 4 (three index levels, one data block). - orders outer: 400 + 10,000 × 4 = 40,400.
That is no better than the worst-case block nested loop, because we probe once for each of 10,000 rows. But suppose a filter first reduces orders to the 100 rows with status = 'pending' (and we read them via the index on status, or they are already in memory). Then: 400 + 100 × 4 = 800, or even less without the scan. Index nested loop is excellent when the outer side is small after filtering and the inner side has an index; it is the typical plan for OLTP queries such as "this user's last 20 orders".
(In practice the upper levels of the B+ tree stay in the buffer pool, so real lookups often cost one or two I/Os, not four.)
4. Sort-merge join
Sort both relations on the join column, then walk through them together like merging two sorted lists, outputting matching pairs. Works only for equality (and some inequality) joins.
orders sorted by customer_id customers sorted by id
[1, 1, 2, 4, 4, 4, 7 ...] [1, 2, 3, 4, 5, 6, 7 ...]
^ advance whichever pointer is smaller; output on equality
- If both inputs are already sorted (for example, read in order from an index): cost = b_r + b_s = 400 + 100 = 500.
- Otherwise add the cost of external merge sort. With M buffer blocks, sort a relation of b blocks by creating ⌈b / M⌉ sorted runs, then merging M − 1 runs at a time. Cost (not counting the final write, since the output goes straight into the merge) = b × (2 × ⌈log_(M−1)(b / M)⌉ + 1).
With M = 12:
- orders: ⌈400 / 12⌉ = 34 runs. Merge passes = ⌈log₁₁ 34⌉ = 2 (one pass handles up to 11 runs, two passes up to 121). Cost = 400 × (2 × 2 + 1) = 2,000.
- customers: ⌈100 / 12⌉ = 9 runs. Merge passes = ⌈log₁₁ 9⌉ = 1. Cost = 100 × (2 × 1 + 1) = 300.
- Merge phase: 400 + 100 = 500.
- Total = 2,000 + 300 + 500 = 2,800.
Sort-merge shines when inputs are already sorted, when the output must be sorted anyway (an ORDER BY on the join key), or when the inputs are huge and many duplicates match.
5. Hash join
Partition both relations with the same hash function on the join column, so matching rows land in matching partitions; then join each pair of partitions in memory.
Build phase: hash the smaller input (customers) into a hash table
Probe phase: for each orders row, hash customer_id and look up matches
If the build side does not fit in memory (Grace hash join):
1. partition customers by h(id) read 100, write 100
2. partition orders by h(customer_id) read 400, write 400
3. for each i: build on customers_i,
probe with orders_i read 100 + 400
- If the build input (customers, 100 blocks) fits in memory: cost = b_r + b_s = 500.
- Otherwise, Grace hash join without recursive partitioning: cost = 3 × (b_r + b_s) = 3 × 500 = 1,500. Each relation is read, written as partitions, and read again.
- Requirement for a single partitioning pass: each build partition must fit in memory. With M = 12, we get M − 1 = 11 partitions of about 100 / 11 ≈ 9.1 blocks, which fit. (The rule of thumb is that M must exceed roughly the square root of the build side's size in blocks.)
Hash join only works for equality conditions, and suffers if the data is skewed (one huge partition for a very common key).
Summary
With M = 12 and no useful filters:
| Algorithm | Cost (block transfers) | Needs | Best when |
|---|---|---|---|
| Nested loop (row-at-a-time) | 1,000,400 | nothing | Tiny inputs, non-equality conditions |
| Block nested loop | 4,100 (customers outer) | nothing | No index, non-equality joins, small inputs |
| Index nested loop | 40,400 (800 with 100 outer rows) | Index on inner join column | Small outer after filtering (OLTP lookups) |
| Sort-merge | 2,800 (500 if pre-sorted) | Equality, sortable keys | Inputs already sorted, sorted output needed |
| Hash join | 1,500 (500 if build fits in memory) | Equality | Large unsorted inputs, analytics |
Interview tip
When asked "which join algorithm would you choose", answer with conditions, not a single name: "Index nested loop when one side is small and the other has an index on the join key; hash join for large equality joins with no useful order; sort-merge when inputs are already sorted or the result must be sorted; nested loop only for tiny inputs or non-equality conditions." Then mention memory: hash and sort spill to disk when their memory budget (PostgreSQL's work_mem) is too small.
MySQL used only nested-loop variants for most of its history; it added hash join in version 8.0.18. PostgreSQL implements all three main families (nested loop, merge join, hash join).
Heuristic optimisation
Before costing anything, optimisers apply equivalence rules of relational algebra that almost always help. The main heuristics:
- Push selections down. Apply
WHEREfilters as early as possible, right above the table scan, so fewer rows flow into joins. Rule used: σ_θ(r ⋈ s) = σ_θ(r) ⋈ s when θ only uses r's columns. - Break up conjunctive selections. σ_(A AND B)(r) = σ_A(σ_B(r)), so each part can be pushed to its own table.
- Push projections down. Drop unneeded columns early so intermediate rows are narrower and more fit per block. Keep only columns needed later (output columns plus join and filter columns).
- Replace a Cartesian product followed by a selection with a join. σ_(r.a = s.b)(r × s) = r ⋈_(r.a = s.b) s, which allows efficient join algorithms.
- Do the most restrictive joins first. Join order matters: joins are commutative and associative, so (r ⋈ s) ⋈ t = r ⋈ (s ⋈ t). Joining the pair that produces the smallest intermediate result first saves work.
Worked example: why pushing down helps
From the earlier query: suppose 1,000 of 5,000 customers are in Pune, and 500 of 10,000 orders are pending, each order has exactly one customer, and the two conditions are independent.
- Join first, then filter. The join produces 10,000 rows (one per order), and only then are filters applied: 10,000 × (500 / 10,000) × (1,000 / 5,000) = 10,000 × 0.05 × 0.2 = 100 rows survive. The join had to process all 10,000 orders and all 5,000 customers.
- Filter first, then join. The join receives 500 orders and 1,000 customers and produces the same 100 rows. The join inputs are 20 times and 5 times smaller respectively, and the hash table holds 1,000 entries instead of 5,000.
Cost-based optimisation
Heuristics cannot decide between, say, a hash join and an index nested loop. For that the optimiser estimates the cost of many alternatives and picks the cheapest.
Statistics
The optimiser's estimates come from statistics stored in the catalog, gathered by ANALYZE (PostgreSQL, and automatically by autovacuum), ANALYZE TABLE (MySQL, plus automatic persistent statistics in InnoDB) or ANALYZE (SQLite):
- Number of rows and pages in each table.
- Per column: number of distinct values (V(A, r)), fraction of NULLs, average width.
- Most common values and their frequencies.
- Histograms describing the distribution of the remaining values (PostgreSQL uses equi-depth histograms: each bucket holds roughly the same number of rows).
- Correlation between physical row order and column order (helps decide between index and sequential scans).
Cardinality estimation
Cardinality here means the number of rows an operator outputs. Selectivity is the fraction of rows that pass a condition. Classic estimates:
| Expression | Estimated rows |
|---|---|
A = constant, no histogram | n_r / V(A, r) (assume values are uniformly distributed) |
A = constant, value in most-common list | n_r × its recorded frequency |
A > constant | n_r × (max − constant) / (max − min), or from the histogram |
C1 AND C2 | n_r × sel(C1) × sel(C2) (assumes independence) |
C1 OR C2 | n_r × (1 − (1 − sel(C1)) × (1 − sel(C2))) |
| r ⋈ s on A, where A is a key of s and a foreign key in r | n_r (each r row matches exactly one s row) |
| r ⋈ s on A, general | n_r × n_s / max(V(A, r), V(A, s)) |
Worked example. orders: 10,000 rows, status has 4 distinct values; customers: 5,000 rows, city has 10 distinct values, no histograms.
- σ status = 'pending' (orders): 10,000 / 4 = 2,500 rows.
- σ city = 'Pune' (customers): 5,000 / 10 = 500 rows.
- Join on customer_id = id, where id is the key of customers: each order matches one customer. Combining with both filters under independence: 10,000 × (1/4) × (1/10) = 250 rows.
If in reality 95 % of orders are 'shipped' and only 1 % are 'pending', the uniform estimate (2,500) is 25 times too high, and the optimiser might choose a sequential scan where an index would be better. This is exactly why databases keep most-common-value lists and histograms.
Why estimates go wrong
- Correlated columns.
city = 'Mumbai' AND state = 'Maharashtra'are not independent; multiplying selectivities underestimates the result badly. PostgreSQL offersCREATE STATISTICSfor multi-column (extended) statistics. - Stale statistics after a bulk load or a big delete. Run
ANALYZE. - Functions and expressions in predicates (
WHERE lower(email) = ...) have no statistics unless you create an expression index. - Errors multiply through joins. A 10x error at each of three joins becomes a 1,000x error at the top.
Plan enumeration
With n tables there are a huge number of join orders (for n = 10, more than 17 billion orderings of bushy trees). Classic optimisers (the System R approach) use dynamic programming: find the best plan for every subset of tables, building larger subsets from smaller ones, and usually restrict the search to left-deep trees (the right child of every join is a base table), which pipeline well. PostgreSQL switches from exhaustive search to a genetic algorithm (GEQO) once a query has 12 or more FROM items by default (geqo_threshold).
Reading EXPLAIN and EXPLAIN ANALYZE
EXPLAIN shows the plan the optimiser chose, with estimates. EXPLAIN ANALYZE actually runs the query and adds real timings and row counts. (Careful: EXPLAIN ANALYZE on an UPDATE or DELETE really changes data; wrap it in BEGIN ... ROLLBACK.)
A PostgreSQL example
The following is PostgreSQL output for orders with 100,000 rows stored in 791 pages and no index on customer_id. The exact numbers depend on your data, version and settings.
EXPLAIN ANALYZE SELECT id, total FROM orders WHERE customer_id = 42;
Seq Scan on orders (cost=0.00..2041.00 rows=20 width=12)
(actual time=0.031..9.842 rows=20 loops=1)
Filter: (customer_id = 42)
Rows Removed by Filter: 99980
Planning Time: 0.095 ms
Execution Time: 9.871 ms
How to read it, piece by piece:
- Node type:
Seq Scanmeans a full sequential scan of the table. - cost=0.00..2041.00: the estimated cost in arbitrary planner units. The first number is the startup cost (before the first row can be returned); the second is the total cost. For a sequential scan with a filter, PostgreSQL's default formula is pages ×
seq_page_cost(1.0) + rows ×cpu_tuple_cost(0.01) + rows ×cpu_operator_cost(0.0025) = 791 + 1,000 + 250 = 2,041. - rows=20: estimated output rows. width=12: estimated average row size in bytes.
- actual time=0.031..9.842: real milliseconds to the first row and to the last row (per loop).
- rows=20 loops=1: real rows returned and how many times this node ran. For a node inside a nested loop, multiply by
loopsto get the total. - Rows Removed by Filter: 99980: the scan read 100,000 rows to return 20. That is the red flag.
Add an index and run it again:
CREATE INDEX idx_orders_customer ON orders (customer_id);
Index Scan using idx_orders_customer on orders
(cost=0.29..80.65 rows=20 width=12)
(actual time=0.024..0.061 rows=20 loops=1)
Index Cond: (customer_id = 42)
Planning Time: 0.120 ms
Execution Time: 0.083 ms
The scan now goes straight to the 20 matching rows: Index Cond instead of Filter, and about 0.08 ms instead of about 10 ms.
A join plan
Hash Join (cost=103.00..2229.13 rows=1000 width=8)
(actual time=1.402..21.774 rows=1000 loops=1)
Hash Cond: (o.customer_id = c.id)
-> Seq Scan on orders o (cost=0.00..2041.00 rows=5000 width=12)
(actual time=0.012..18.905 rows=5000 loops=1)
Filter: (status = 'pending'::text)
Rows Removed by Filter: 95000
-> Hash (cost=90.50..90.50 rows=1000 width=4)
(actual time=1.371..1.372 rows=1000 loops=1)
Buckets: 1024 Batches: 1 Memory Usage: 44kB
-> Seq Scan on customers c (cost=0.00..90.50 rows=1000 width=4)
(actual time=0.008..1.150 rows=1000 loops=1)
Filter: (city = 'Pune'::text)
Rows Removed by Filter: 4000
Read plans from the innermost (most indented) nodes outwards: the customers scan feeds the Hash node (the build side, which fits in memory: Batches: 1), and the orders scan is the probe side. Batches greater than 1 means the hash table spilled to disk because work_mem was too small.
What to look for in any plan:
| Symptom | Likely cause | Fix |
|---|---|---|
| Seq Scan with a large "Rows Removed by Filter" on a selective condition | Missing index | Add an index on the filter column(s) |
| Estimated rows far from actual rows (10x or more) | Stale or missing statistics, correlated columns | ANALYZE; extended statistics |
Nested Loop with a high loops count on the inner side | Bad row estimate for the outer side | Fix estimates; check indexes on the join column |
Sort Method: external merge Disk: ... | Sort exceeded work_mem | Index matching ORDER BY, or more work_mem for that query |
Hash Batches greater than 1 | Hash table did not fit in memory | More work_mem, or reduce rows earlier |
Other databases have the same idea with different output: MySQL has EXPLAIN (a table with type, key, rows columns, where type = ALL means a full scan) and EXPLAIN ANALYZE (8.0.18+). SQLite has EXPLAIN QUERY PLAN. In SQLite 3.51 with the 100,000-row table above:
sqlite> EXPLAIN QUERY PLAN SELECT id, total FROM orders WHERE customer_id = 42;
QUERY PLAN
`--SEARCH orders USING INDEX idx_orders_customer (customer_id=?)
SCAN means a full table scan; SEARCH ... USING INDEX means an index lookup.
Common slow-query fixes
1. The N+1 query problem
An application loads a list, then runs one more query for each item:
orders = db.query("SELECT id, customer_id FROM orders LIMIT 50") # 1 query
for o in orders:
o.customer = db.query("SELECT name FROM customers WHERE id = ?", # 50 queries
o.customer_id)
That is 51 round trips. Each is fast, but 50 network round trips of 1 ms each add 50 ms, and it gets worse with page size. ORMs cause this silently through lazy loading. Fix it with a join or a single batched query:
SELECT o.id, c.name
FROM orders o
JOIN customers c ON c.id = o.customer_id
LIMIT 50;
-- or, two queries in total:
SELECT name, id FROM customers WHERE id IN (3, 17, 42, ...);
In ORMs this is "eager loading": select_related / prefetch_related in Django, includes in Rails, JOIN FETCH or entity graphs in JPA.
2. Missing or unusable index
- Add indexes on columns used in
WHERE,JOINandORDER BYof frequent queries. - For multi-column conditions, use a composite index with equality columns first, then the range or sort column:
(status, created_at)servesWHERE status = ? ORDER BY created_at DESC. - An index is not used when the column is wrapped in a function or cast:
WHERE lower(email) = 'a@b.com'orWHERE DATE(created_at) = '2026-10-10'. Rewrite to a range (created_at >= '2026-10-10' AND created_at < '2026-10-11') or create an expression index (CREATE INDEX ON users (lower(email)), supported by PostgreSQL and SQLite; MySQL 8.0.13+ supports functional key parts). - A leading wildcard (
LIKE '%phone') cannot use a normal B+ tree index. Use a full-text or trigram index. - Implicit type conversions (comparing a
VARCHARcolumn with a number in MySQL) can also stop index use.
3. SELECT *
SELECT * fetches every column, even ones you do not need:
- More bytes read from disk and sent over the network (large
TEXTandJSONcolumns hurt most). - It prevents index-only scans (covering indexes): if the query needs only
idandcreated_at, an index on(status, created_at)that includes them can answer it without touching the table. SQLite shows this asUSING COVERING INDEX, PostgreSQL asIndex Only Scan. - Code breaks or silently changes when someone adds a column.
List the columns you need.
4. OFFSET pagination versus keyset pagination
OFFSET pagination looks innocent:
SELECT id, created_at FROM orders
WHERE status = 'pending'
ORDER BY created_at DESC, id DESC
LIMIT 20 OFFSET 100000;
But the database must produce and throw away the first 100,000 rows to return 20. Page 5,001 is about 5,000 times more work than page 1. Results also shift when rows are inserted between page loads (you see duplicates or miss rows).
Keyset pagination (also called seek or cursor pagination) remembers the last row of the previous page and asks for rows after it:
-- page 1
SELECT id, created_at FROM orders
WHERE status = 'pending'
ORDER BY created_at DESC, id DESC
LIMIT 20;
-- next page: pass the last row's (created_at, id) as the cursor
SELECT id, created_at FROM orders
WHERE status = 'pending'
AND (created_at, id) < ('2026-09-18', 99380)
ORDER BY created_at DESC, id DESC
LIMIT 20;
With an index on (status, created_at, id) (or (status, created_at) in databases where the primary key is implicitly part of the index, as in SQLite and InnoDB), every page is a short index range scan: the cost is the same for page 1 and page 5,000. The id tie-breaker makes the order unique, so no row is skipped or repeated when several rows share a timestamp. Row-value comparisons like (created_at, id) < (...) work in PostgreSQL, MySQL and SQLite 3.15+. The trade-off: you cannot jump to "page 37" directly, which is fine for infinite scroll and "next" buttons.
Tested in SQLite: on 100,000 orders, the keyset query for page 2 returned exactly the same rows as LIMIT 3 OFFSET 3, and the plan was SEARCH orders USING COVERING INDEX idx_orders_status_created (status=? AND created_at<?).
5. Other frequent culprits
- Unbounded queries: always
LIMITlist endpoints. COUNT(*)on huge tables for every page view: cache it, or show an estimate.ORacross different columns that defeats indexes: rewrite asUNION ALLof two indexed queries when needed.NOT INwith a subquery that may contain NULLs: it returns no rows at all; useNOT EXISTS.- Too many indexes on a write-heavy table: every insert updates every index. Remove unused ones.
- Lock waits that look like slow queries: check for long transactions holding locks (see concurrency control).
Interview tip
For "how would you debug a slow query?", give a process: (1) reproduce it and get the exact SQL and parameters, (2) run EXPLAIN ANALYZE, (3) find the node where most time goes and compare estimated with actual rows, (4) fix the cause: index, rewrite, statistics, or less data, (5) measure again. Mention the slow query log (MySQL) or pg_stat_statements (PostgreSQL) for finding the worst queries in the first place.
Common mistake
Adding an index for every column "to be safe" is not free. Each index slows down inserts, updates and deletes, uses disk and memory, and an index on a low-selectivity column (such as a boolean) is often never used. Index for real query patterns, and check that the plan uses it.
Interview questions
Q1. What are the steps of query processing?
Parsing checks syntax and resolves names against the catalog; rewriting expands views and simplifies the query into a logical plan; optimisation enumerates physical plans and picks the cheapest by estimated cost; execution runs the plan, usually as a tree of iterators pulling rows from their children. Plans can be cached for prepared statements.
Q2. Why push selections down in a query tree?
Filtering early shrinks the inputs of later, more expensive operators such as joins and sorts. Smaller inputs mean fewer I/Os, smaller hash tables and less memory. It is almost always beneficial, which is why optimisers apply it as a heuristic before cost-based search.
Q3. Compare nested loop, sort-merge and hash joins.
Nested loop compares each outer row with the inner input; it works for any condition and is fast with an index on the inner side and a small outer side. Sort-merge sorts both inputs on the key and merges them; it is good when inputs are already sorted. Hash join builds a hash table on the smaller input and probes with the larger; it is usually the best for large equality joins. Hash and sort-merge need equality (or ordering) conditions.
Q4. What is the cost of a block nested loop join?
With M buffer blocks, it is ⌈b_outer / (M − 2)⌉ × b_inner + b_outer block transfers; in the worst case with one outer block at a time it is b_outer × b_inner + b_outer. Choosing the smaller relation as the outer side and giving it more memory reduces the cost.
Q5. What is the cost of a Grace hash join?
About 3 × (b_r + b_s) block transfers when no recursive partitioning is needed: both relations are read once, written as partitions, and read again for the build and probe. If the build side fits in memory, it is just b_r + b_s.
Q6. Why might the optimiser not use my index?
The condition may not be selective enough, so a sequential scan is cheaper than many random reads. Or the index cannot be used: the column is wrapped in a function, there is a type mismatch, a leading wildcard in LIKE, or the composite index's leading column is not in the query. Statistics may also be stale, making the optimiser misjudge selectivity.
Q7. What statistics does a cost-based optimiser use?
Table row and page counts, and per-column distinct counts, NULL fractions, most-common values with frequencies, histograms and physical correlation. They are collected by ANALYZE commands or automatically. From them the optimiser estimates selectivity and cardinality for each operator.
Q8. What is cardinality estimation and why is it hard?
It is predicting how many rows each operator will produce. It is hard because optimisers assume uniform distributions and independent columns, which real data violates, and errors multiply across joins. A bad estimate can make the optimiser pick a nested loop that runs millions of times instead of a hash join.
Q9. What is the difference between EXPLAIN and EXPLAIN ANALYZE?
EXPLAIN shows the chosen plan with estimated costs and row counts without running the query. EXPLAIN ANALYZE executes it and adds actual times, row counts and loop counts, so you can compare estimates with reality. Because it executes the statement, wrap data-changing statements in a transaction and roll back.
Q10. What is the N+1 query problem?
Loading a list with one query and then running one extra query per item, usually through lazy-loaded ORM relationships. It causes many round trips, which adds latency linearly with the number of items. Fix it with a join or a single batched IN query (eager loading).
Q11. Why is OFFSET pagination slow for deep pages?
The database must generate and discard all the skipped rows, so the cost grows linearly with the offset. Keyset pagination instead filters by the last seen sort key, for example WHERE (created_at, id) < (?, ?), which is an index range scan with constant cost per page and stable results under concurrent inserts.
Q12. What is a covering index?
An index that contains every column a query needs, so the database can answer from the index alone without reading the table rows (an index-only scan). It saves random I/O. Avoiding SELECT * is what makes covering indexes possible.
Q13. How does the optimiser choose a join order?
It uses dynamic programming over subsets of tables, keeping the cheapest plan for each subset and combining them, typically restricted to left-deep trees. For queries with many tables, it falls back to heuristics or randomised search (PostgreSQL's genetic optimiser). Join order matters because it controls the sizes of intermediate results.
Q14. What is pipelining?
Passing rows from one operator directly to the next as they are produced, without writing intermediate results to disk. It reduces I/O and gives the first rows sooner. Blocking operators such as sort, aggregation and the build side of a hash join must consume their full input before producing output.
Key takeaways
- A query is parsed, rewritten, optimised and executed; execution is usually a pipeline of iterators.
- Relational algebra trees represent plans; equivalent trees can differ in cost by orders of magnitude.
- A secondary index helps only selective conditions; for large fractions of a table a sequential scan wins.
- Join costs for our example (400 and 100 blocks, M = 12): nested loop 1,000,400; block nested loop 4,100; sort-merge 2,800; hash join 1,500; index nested loop 800 when the outer side is filtered to 100 rows.
- Heuristics: push selections and projections down, turn products into joins, do restrictive joins first.
- Cost-based optimisation relies on statistics; stale stats and correlated columns cause bad plans.
- In
EXPLAIN ANALYZE, read innermost nodes first and compare estimated and actual rows. - Everyday fixes: remove N+1 queries, add the right composite index, avoid functions on indexed columns, select only needed columns, and use keyset pagination.
Next lesson
Continue with NoSQL and distributed databases.

