Join Algorithms - Nested Loop vs Index Nested Loop vs Hash Join vs Merge Join
Nested loop vs index nested loop vs hash join vs merge join: runnable work-count model, textbook I/O cost formulas and Postgres-style batch math (runnable), real PG 17 plans forcing each algorithm (incl. Memoize) and a hash join spilling to Batches: 8 under small work_mem; when each wins table, the planned-10-got-1M nested loop failure, MySQL (no merge join, hash join 8.0.18+), SQLite automatic indexes, DuckDB range joins.
- 1Gist
- 2Maps
- 3Q&A
- 4Sandbox
Voice readout needs Web Speech Synthesis in this browser.
Question ladder
L1
Which join algorithm handles any join condition?
Answer
Nested loop. Hash and merge joins need an equality (merge also supports merge-able operators).
L2
What is the cost shape of an index nested loop?
Answer
Outer rows times one index probe, roughly |R| x log|S|.
L3
Which side should a hash join build on?
Answer
The smaller side by estimated bytes, because memory and batches scale with the build side.
L4
What does Batches: 8 on a Hash node mean?
Answer
The build side did not fit, so both inputs were partitioned and seven batches went through temp files.
L5
When does merge join shine?
Answer
When both inputs already arrive sorted on the key, from indexes or an earlier merge, so the join is one pass with tiny memory.
L6
What is Memoize?
Answer
A cache between a nested loop and its inner side, keyed on the outer join key, that turns repeated probes into hits.
L7
Why not set enable_nestloop = off globally?
Answer
It fixes one query and breaks hundreds of point lookups that need index nested loops.
Failure modes
Nested loop over a huge outer input
The planner expected a few outer rows, got millions, and repeated the inner probe millions of times.
Hash join spills
The build side was underestimated or work_mem is small, so the hash table splits into batches and every spilled row is written and read again.
Type mismatch on the join key
A cast on the inner key hides its index, turning an index nested loop into a full-scan hash join.
Misconceptions
Hash join is always fastest.
For 20 lookups into an indexed table, building a hash over the whole inner table is far more work than 20 probes.
Memoize is a fourth join algorithm.
It is a cache that speeds up nested loops with repeated keys.
Rows on the inner side of a nested loop are totals.
They are per loop. Multiply by loops.
Interviewer traps
Fixing a bad nested loop by disabling nested loops.
Fix the outer row estimate instead so the planner can see that a hash join is cheaper.
Ignoring join fan-out.
A many-to-many join can return far more rows than either input. Estimate the output, not only the inputs.
Design scenario
Same prompt for every reader.
Requirements
The API endpoint returns one user's last 50 events in under 20 ms; the nightly report joins everything within an hour.
Failure assumptions
- Statistics on events.user_id lag behind the data.
- Several reports run at the same time.
- Some users have millions of events.
Constraints
- PostgreSQL 17 with 64 GB RAM.
- Indexes on users.id and events.user_id only.
Prompt
An events table with 300 million rows is joined to a 5-million-row users table for a nightly report and for a per-user API endpoint. Choose join strategies and memory settings for both.
API
Which queries does each path send, and which join algorithm should each get?
Data
Which indexes and statistics make the API path an index nested loop and keep the report a hash join?
Architecture
How do you size work_mem for the report without risking memory for the API fleet?
Overview
Every join in a plan is executed by one of a handful of physical algorithms, and the planner picks one per join. Nested loop re-runs the inner side for each outer row and is unbeatable when the outer side is tiny and the inner side has an index (index nested loop). Hash join builds a hash table on the smaller input and streams the larger one through it: the default winner for large equi-joins, but it needs memory and spills to disk in batches when the build side exceeds work_mem. Merge join walks two inputs sorted on the join key in lockstep: great when both sides already arrive sorted (from indexes or an earlier sort), and the only scalable choice for some big-big joins when memory is tight. The planner's choice is only as good as its row estimates: a nested loop planned for 10 outer rows that actually receives 2 million is the single most common "query went from 50 ms to 20 minutes" story.
The four algorithms, side by side
Decisions
- 1
Step 1: Join R and S on R.k = S.k
- nextStep 2: Is the condition an equality?
- ?
Step 2: Is the condition an equality?
- no: range, LIKE, inequalityStep 3a: Nested Loop - only general option, ideally with an index range probe on the inner side
- yesStep 3b: Is the outer side small and the inner side indexed on k?
- 3
Step 3a: Nested Loop - only general option, ideally with an index range probe on the inner side
- ?
Step 3b: Is the outer side small and the inner side indexed on k?
- yesStep 4a: Index Nested Loop - cost ~ outer rows x index probe
- noStep 4b: Do both inputs already arrive sorted on k?
- 5
Step 4a: Index Nested Loop - cost ~ outer rows x index probe
- ?
Step 4b: Do both inputs already arrive sorted on k?
- yesStep 5a: Merge Join - one pass over each input
- noStep 5b: Hash Join - build on smaller input, probe with larger
- 7
Step 5a: Merge Join - one pass over each input
- 8
Step 5b: Hash Join - build on smaller input, probe with larger
- nextStep 6: Build side bigger than work_mem x hash_mem_multiplier?
- ?
Step 6: Build side bigger than work_mem x hash_mem_multiplier?
- yesStep 7: Batches > 1 - partitions written to temp files, read back
- noStep 7: Single batch, all in memory
- 10
Step 7: Batches > 1 - partitions written to temp files, read back
- 11
Step 7: Single batch, all in memory
Lesson map
Join Algorithms - Nested Loop vs Index Nested Loop vs Hash Join vs Merge Join
Nested loop vs index nested loop vs hash join vs merge join: runnable work-count model, textbook I/O cost formulas and Postgres-style batch math (runnable), real PG 17 plans forcing each algorithm (incl. Memoize) and a hash join spilling to Batches: 8 under small work_mem; when each wins table, the planned-10-got-1M nested loop failure, MySQL (no merge join, hash join 8.0.18+), SQLite automatic indexes, DuckDB range joins.
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 q["Step 1: Join R and S on R.k = S.k"] d["Step 2: Is the condition an equality?"] nl["Step 3a: Nested Loop - only general option, ideally with an index range probe on the inner side"] sz["Step 3b: Is the outer side small and the inner side indexed on k?"] inl["Step 4a: Index Nested Loop - cost ~ outer rows x index probe"] so["Step 4b: Do both inputs already arrive sorted on k?"] mj["Step 5a: Merge Join - one pass over each input"] hj["Step 5b: Hash Join - build on smaller input, probe with larger"] sp["Step 6: Build side bigger than work_mem x hash_mem_multiplier?"] bt["Step 7: Batches > 1 - partitions written to temp files, read back"] mem["Step 7: Single batch, all in memory"] q -->|continues| d d -->|no: range, LIKE, inequality| nl d -->|yes| sz sz -->|yes| inl sz -->|no| so so -->|yes| mj so -->|no| hj hj -->|continues| sp sp -->|yes| bt sp -->|no| mem
Simple nested loop. For every row in the outer input, scan the entire inner input and test the condition. Cost grows as |R| x |S|. Nobody wants this on large inputs, but it is the only algorithm that handles any join condition, and it is perfectly fine when one side has a handful of rows.
Index nested loop. Same loop, but the inner side is an index lookup on the join key. Cost is |R| x (index descent + matching rows), roughly |R| x log|S|. This is the OLTP workhorse: "fetch these 20 orders and their customers" is 20 primary-key probes. It also pipelines perfectly and returns first rows immediately, which is why it pairs so well with LIMIT.
Hash join. Two phases. Build: read the smaller input (by estimated size) and insert each row into an in-memory hash table keyed on the join columns. Probe: stream the larger input and look up each row's key. Cost is roughly |R| + |S| with a memory requirement proportional to the build side. Only works for equality conditions. When the build side does not fit, PostgreSQL splits both inputs into batches by hash value, keeps one batch in memory and writes the rest to temp files (a variant of the classic Grace hash join), so every spilled row is written and read once more.
Merge join (sort-merge). Both inputs sorted on the key; advance two cursors, emitting matches when keys are equal (with care for duplicates on both sides). If the inputs already arrive sorted (an index scan on the key, or the output of an earlier merge join or sort), the join itself is a single pass with tiny memory. If they must be sorted first, the sort cost dominates and may spill (see the sorting page). Merge join also supports some range-ish conditions better than hash joins, and produces sorted output that a later ORDER BY or GROUP BY can reuse.
When each algorithm wins
| Situation | Best algorithm | Why | What happens if the planner picks something else |
|---|---|---|---|
| Few outer rows, inner indexed on key (OLTP lookups) | Index nested loop | Cost proportional to outer rows; first row immediately | Hash join reads the whole inner table to build a hash for 20 lookups |
| Two large inputs, equality, no useful order | Hash join | Linear in both inputs | Nested loop is quadratic; merge join pays two big sorts |
| Both inputs already sorted on key (indexes, previous merge) | Merge join | One pass, tiny memory, sorted output reused by ORDER BY | Hash join discards the order and may need a later sort |
| Build side larger than memory, many concurrent queries | Merge join or hash join with batches | Predictable memory | Single-batch hash estimate on bad stats can exceed memory; batches double the I/O |
| Non-equality condition (ranges, geometric, LIKE) | Nested loop with index range probe | Only general algorithm | Hash and merge are not applicable |
| Outer estimate wrong: planned 10 rows, gets 1M | Hash join would have been right | Nested loop repeats the inner probe 1M times | The classic regression: minutes instead of milliseconds |
Hash join or nested loop for a big equi-join?
Prefer
Hash join
Build a hash table on the smaller input once, then stream the larger input through it.
- 21,000 work units for 20,000 orders x 1,000 customers.
- Real PostgreSQL chose it freely for 300,000 events x 50,000 users.
- It needs memory: 256kB of work_mem gave Batches: 8.
Alternative
Nested loop
Re-run the inner side for every outer row.
- The naive version did 20,000,000 comparisons.
- The index version did 200,000 work units, about 10x the hash join.
- It wins only when the outer side is tiny or the condition is a range.
How the planner narrows the join algorithm
Diagram 1 condensed.
- 1
Check the condition
No equality means nested loop, ideally with an index range probe. - 2
Small outer, indexed inner
Index nested loop: cost tracks outer rows. - 3
Both sides sorted
Merge join: one pass over each input. - 4
Otherwise hash
Build on the smaller input and probe with the larger. - 5
Build side too big
Batches above 1 write partitions to temp files.
Measuring the work, algorithm by algorithm
The runnable model joins 20,000 orders to 1,000 customers four ways and counts comparisons or probes. The absolute numbers are a proxy, but the ratios are the point.
"""Four join algorithms on the same data, counting the work each does.
R = orders (outer/probe side), S = customers (inner/build side). Join on customer_id = id.
"Comparisons" and "probes" are a rough proxy for CPU work; sizes are example values.
"""
import bisect
from collections import defaultdict
orders = [(i, (i * 37) % 2000) for i in range(20_000)] # (order_id, customer_id)
customers = [(c, f"cust{c}") for c in range(0, 2000, 2)] # only even ids exist -> half the orders match
def nested_loop(R, S):
work, out = 0, []
for o in R: # for every outer row...
for c in S: # ...scan the whole inner input
work += 1
if o[1] == c[0]: out.append((o[0], c[1]))
return out, work
def index_nested_loop(R, S):
keys = [c[0] for c in S] # pretend this sorted array is a B-tree on customers.id
work, out = 0, []
for o in R:
j = bisect.bisect_left(keys, o[1]) # O(log |S|) descent per outer row
work += max(1, len(keys).bit_length())
if j < len(keys) and keys[j] == o[1]: out.append((o[0], S[j][1]))
return out, work
def hash_join(R, S):
table = defaultdict(list)
for c in S: table[c[0]].append(c) # build: hash the smaller input once
work, out = len(S), []
for o in R: # probe: one hash lookup per outer row
work += 1
for c in table.get(o[1], ()): out.append((o[0], c[1]))
return out, work
def merge_join(R, S):
Rs, Ss = sorted(R, key=lambda o: o[1]), sorted(S, key=lambda c: c[0]) # both sides sorted on the key
i = j = work = 0; out = []
while i < len(Rs) and j < len(Ss):
work += 1
if Rs[i][1] < Ss[j][0]: i += 1
elif Rs[i][1] > Ss[j][0]: j += 1
else:
k = j # emit all inner rows with this key
while k < len(Ss) and Ss[k][0] == Rs[i][1]:
out.append((Rs[i][0], Ss[k][1])); k += 1
i += 1
return out, work
results = {}
for name, fn in [("nested loop", nested_loop), ("index nested loop", index_nested_loop),
("hash join", hash_join), ("merge join (excl. sort)", merge_join)]:
out, work = fn(orders, customers)
results[name] = sorted(out)
print(f"{name:24s} rows={len(out):6d} work units={work:>12,d}")
print("all four agree:", len({tuple(v) for v in results.values()}) == 1)
# Non-equi join (range predicate): hash join cannot help, the key is not an equality
band = [(o, c) for o in orders[:200] for c in customers if abs(o[1] - c[0]) <= 2]
print("band join (|o.cust - c.id| <= 2) rows:", len(band), "-> only nested loop (or a range index probe) applies")Output:
nested loop rows= 10000 work units= 20,000,000
index nested loop rows= 10000 work units= 200,000
hash join rows= 10000 work units= 21,000
merge join (excl. sort) rows= 10000 work units= 20,990
all four agree: True
band join (|o.cust - c.id| <= 2) rows: 498 -> only nested loop (or a range index probe) appliesTwenty million comparisons for the naive nested loop versus about twenty thousand for hash and merge joins: three orders of magnitude, on tiny inputs. Index nested loop sits in between: excellent when the outer side is small, too expensive when the outer side is the 20,000-row table and every probe walks a B-tree. And note the last line: for a band condition (abs(a - b) <= 2) there is no hash key, so only a nested loop (ideally with an index range probe on the inner side) can run it.
The cost formulas interviewers expect
Textbook I/O costs (CMU 15-445 style) with R = M pages, S = N pages, B buffer pages, m rows in R:
| Algorithm | I/O cost (pages) | Memory | Condition types | Output order |
|---|---|---|---|---|
| Simple nested loop | M + m x N | ~3 pages | Any | Outer order |
| Block nested loop | M + ceil(M / (B-2)) x N | B pages | Any | Outer order (by chunk) |
| Index nested loop | M + m x (index probe cost) | Small | Equality or range on indexed column | Outer order |
| Sort-merge | sort(M) + sort(N) + M + N | B pages for sorts | Equality (and merge-able operators) | Sorted on key |
| Hash join (fits) | M + N | Build side x overhead | Equality only | None |
| Grace hash join (spills) | 3 x (M + N) | B pages | Equality only | None |
// Textbook I/O cost formulas (CMU 15-445 style) for joining R (M pages) with S (N pages)
// using B buffer pages. All sizes are example values; real optimizers add CPU terms.
type Costs = Record<string, number>;
function joinCosts(M: number, N: number, B: number, mRows: number, probeIo: number): Costs {
const sortCost = (P: number): number => {
// external merge sort: 2P per pass, passes = 1 + ceil(log_{B-1}(ceil(P/B)))
const runs = Math.ceil(P / B);
const passes = runs <= 1 ? 1 : 1 + Math.ceil(Math.log(runs) / Math.log(B - 1));
return 2 * P * passes;
};
const hashFits = Math.min(M, N) <= B - 2; // build side fits in memory -> one pass, no spill
return {
"simple nested loop": M + mRows * N, // re-read S once per R *row*
"block nested loop": M + Math.ceil(M / (B - 2)) * N, // re-read S once per R *chunk*
"index nested loop": M + mRows * probeIo, // one index probe per R row
"sort-merge (sort both)": sortCost(M) + sortCost(N) + M + N,
"hash join": hashFits ? M + N : 3 * (M + N), // grace hash: partition (write+read) then join
};
}
// Postgres-style batch count: inner bytes / work_mem, rounded up to a power of two.
function hashBatches(innerRows: number, rowBytes: number, workMemKb: number, hashMemMultiplier = 2): number {
const bytes = innerRows * rowBytes;
const budget = workMemKb * 1024 * hashMemMultiplier;
let b = 1;
while (bytes / b > budget) b *= 2;
return b;
}
const scenarios = [
{ name: "big x big, small memory", M: 1000, N: 500, B: 100, mRows: 100_000, probeIo: 3 },
{ name: "big x big, plenty memory", M: 1000, N: 500, B: 600, mRows: 100_000, probeIo: 3 },
{ name: "tiny outer x indexed big", M: 1, N: 50_000, B: 100, mRows: 20, probeIo: 3 },
];
for (const s of scenarios) {
const c = joinCosts(s.M, s.N, s.B, s.mRows, s.probeIo);
const best = Object.entries(c).sort((a, b) => a[1] - b[1])[0];
console.log(`${s.name}: ` + Object.entries(c).map(([k, v]) => `${k}=${v.toLocaleString("en-US")}`).join(", "));
console.log(` -> cheapest: ${best[0]}`);
}
// Example: 50k-row build side at ~48 bytes/row vs work_mem settings
for (const wm of [256, 1024, 4096]) {
console.log(`hash build 50,000 rows x 48B with work_mem=${wm}kB -> batches=${hashBatches(50_000, 48, wm)}`);
}Output:
big x big, small memory: simple nested loop=50,001,000, block nested loop=6,500, index nested loop=301,000, sort-merge (sort both)=7,500, hash join=4,500
-> cheapest: hash join
big x big, plenty memory: simple nested loop=50,001,000, block nested loop=2,000, index nested loop=301,000, sort-merge (sort both)=6,500, hash join=1,500
-> cheapest: hash join
tiny outer x indexed big: simple nested loop=1,000,001, block nested loop=50,001, index nested loop=61, sort-merge (sort both)=350,003, hash join=50,001
-> cheapest: index nested loop
hash build 50,000 rows x 48B with work_mem=256kB -> batches=8
hash build 50,000 rows x 48B with work_mem=1024kB -> batches=2
hash build 50,000 rows x 48B with work_mem=4096kB -> batches=1Expectedbig x big, small memory: simple nested loop=50,001,000, block nested loop=6,500, index nested loop=301,000, sort-merge (sort both)=7,500, hash join=4,500 -> cheapest: hash join big x big, plenty memory: simple nested loop=50,001,000, block nested loop=2,000, index nested loop=301,000, sort-merge (sort both)=6,500, hash join=1,500 -> cheapest: hash join tiny outer x indexed big: simple nested loop=1,000,001, block nested loop=50,001, index nested loop=61, sort-merge (sort both)=350,003, hash join=50,001 -> cheapest: index nested loop hash build 50,000 rows x 48B with work_mem=256kB -> batches=8 hash build 50,000 rows x 48B with work_mem=1024kB -> batches=2 hash build 50,000 rows x 48B with work_mem=4096kB -> batches=1
Press Run. Snippets must be self-contained — no network, files, or native modules.
The last three lines are a simplified version of PostgreSQL's batch math: build bytes divided by work_mem x hash_mem_multiplier, rounded up to a power of two. The 256 kB case predicts 8 batches, and the real PostgreSQL run below reports Batches: 8 for a similar build side.
Real PostgreSQL: one join, three algorithms, and a spill
The session below joins 300,000 events to 50,000 users, lets the planner choose, then disables hash join and merge join in turn (diagnostic switches: they raise the cost of that path so the planner picks the next best one), and finally runs the hash join with a tiny and a comfortable work_mem.
-- Page 2: one join, three physical algorithms, plus a hash join that spills when work_mem is small.
\pset footer off
SET client_min_messages = warning;
DROP SCHEMA IF EXISTS p2 CASCADE;
CREATE SCHEMA p2;
SET search_path = p2;
SET max_parallel_workers_per_gather = 0;
CREATE TABLE users AS SELECT g AS id, 'user_' || g AS name FROM generate_series(1, 50000) g;
CREATE TABLE events AS
SELECT g AS id, 1 + (g * 31) % 50000 AS user_id, g % 7 AS kind
FROM generate_series(1, 300000) g;
ALTER TABLE users ADD PRIMARY KEY (id);
CREATE INDEX events_user_idx ON events (user_id);
ANALYZE users; ANALYZE events;
-- 1) Planner's free choice for a big equi-join: Hash Join (build on the smaller input)
EXPLAIN (COSTS OFF) SELECT count(*) FROM events e JOIN users u ON u.id = e.user_id;
-- 2) Disable hash join -> planner falls back to Merge Join (both sides delivered sorted by index)
SET enable_hashjoin = off;
EXPLAIN (COSTS OFF) SELECT count(*) FROM events e JOIN users u ON u.id = e.user_id;
-- 3) Disable merge join too -> Nested Loop with an inner index probe
SET enable_mergejoin = off;
EXPLAIN (COSTS OFF) SELECT count(*) FROM events e JOIN users u ON u.id = e.user_id;
RESET enable_hashjoin; RESET enable_mergejoin;
-- 4) Hash join spill: same query, build side forced larger than work_mem -> Batches > 1 (temp files)
SET work_mem = '256kB';
EXPLAIN (ANALYZE, COSTS OFF, TIMING OFF, SUMMARY OFF)
SELECT count(*) FROM events e JOIN users u ON u.id = e.user_id;
SET work_mem = '16MB';
EXPLAIN (ANALYZE, COSTS OFF, TIMING OFF, SUMMARY OFF)
SELECT count(*) FROM events e JOIN users u ON u.id = e.user_id;
-- 5) Non-equi join (range condition): hash and merge cannot be used, only Nested Loop
EXPLAIN (COSTS OFF)
SELECT count(*) FROM users a JOIN users b ON b.id BETWEEN a.id AND a.id + 2 WHERE a.id < 100;Real output (PostgreSQL 17.11, local sandbox, psql; setup DDL echoes omitted):
SET max_parallel_workers_per_gather = 0;
EXPLAIN (COSTS OFF) SELECT count(*) FROM events e JOIN users u ON u.id = e.user_id;
QUERY PLAN
---------------------------------------
Aggregate
-> Hash Join
Hash Cond: (e.user_id = u.id)
-> Seq Scan on events e
-> Hash
-> Seq Scan on users u
SET enable_hashjoin = off;
EXPLAIN (COSTS OFF) SELECT count(*) FROM events e JOIN users u ON u.id = e.user_id;
QUERY PLAN
---------------------------------------------------------------
Aggregate
-> Merge Join
Merge Cond: (e.user_id = u.id)
-> Index Only Scan using events_user_idx on events e
-> Index Only Scan using users_pkey on users u
SET enable_mergejoin = off;
EXPLAIN (COSTS OFF) SELECT count(*) FROM events e JOIN users u ON u.id = e.user_id;
QUERY PLAN
---------------------------------------------------------------
Aggregate
-> Nested Loop
-> Seq Scan on events e
-> Memoize
Cache Key: e.user_id
Cache Mode: logical
-> Index Only Scan using users_pkey on users u
Index Cond: (id = e.user_id)
RESET enable_hashjoin;
RESET enable_mergejoin;
SET work_mem = '256kB';
EXPLAIN (ANALYZE, COSTS OFF, TIMING OFF, SUMMARY OFF)
SELECT count(*) FROM events e JOIN users u ON u.id = e.user_id;
QUERY PLAN
-------------------------------------------------------------------
Aggregate (actual rows=1 loops=1)
-> Hash Join (actual rows=300000 loops=1)
Hash Cond: (e.user_id = u.id)
-> Seq Scan on events e (actual rows=300000 loops=1)
-> Hash (actual rows=50000 loops=1)
Buckets: 16384 Batches: 8 Memory Usage: 351kB
-> Seq Scan on users u (actual rows=50000 loops=1)
SET work_mem = '16MB';
EXPLAIN (ANALYZE, COSTS OFF, TIMING OFF, SUMMARY OFF)
SELECT count(*) FROM events e JOIN users u ON u.id = e.user_id;
QUERY PLAN
-------------------------------------------------------------------
Aggregate (actual rows=1 loops=1)
-> Hash Join (actual rows=300000 loops=1)
Hash Cond: (e.user_id = u.id)
-> Seq Scan on events e (actual rows=300000 loops=1)
-> Hash (actual rows=50000 loops=1)
Buckets: 65536 Batches: 1 Memory Usage: 2270kB
-> Seq Scan on users u (actual rows=50000 loops=1)
EXPLAIN (COSTS OFF)
SELECT count(*) FROM users a JOIN users b ON b.id BETWEEN a.id AND a.id + 2 WHERE a.id < 100;
QUERY PLAN
-----------------------------------------------------------------
Aggregate
-> Nested Loop
-> Index Only Scan using users_pkey on users a
Index Cond: (id < 100)
-> Index Only Scan using users_pkey on users b
Index Cond: ((id >= a.id) AND (id <= (a.id + 2)))What the real output teaches:
- Free choice is Hash Join, building on
users(the smaller input) and probing with all 300,000 events. Both sides are Seq Scans because every row participates. - Without hash join, Merge Join appears, fed by two Index Only Scans that already deliver rows sorted by the join key: no explicit Sort node needed. If those indexes did not exist, you would see two Sort nodes under the merge join instead.
- Without merge join, Nested Loop with a Memoize node (PostgreSQL 14+): because many events share the same
user_id, the planner caches inner-side results per key, turning repeated probes into cache hits. Memoize is a nested-loop accelerator, not a separate join algorithm. - Spill: with
work_mem = 256kB, the Hash node reportsBatches: 8(the build side was partitioned and seven batches went through temp files). With 16 MB it reportsBatches: 1 Memory Usage: 2270kB. Same plan shape, different physical behaviour, and onlyEXPLAIN ANALYZEshows it. - Range condition:
b.id BETWEEN a.id AND a.id + 2has no equality, so the plan is a Nested Loop with an index range scan on the inner side, exactly as the decision diagram predicts.
The "planner picked wrong" failure, concretely
The most damaging misplan is a nested loop over an outer input that is much bigger than estimated. The planner estimates outer rows = 1 (for example, because of a correlated-columns underestimate or a stale histogram that does not know about yesterday's data), so a nested loop with an inner index probe looks cheapest. In reality the outer side produces 500,000 rows, and the inner probe runs 500,000 times, each a few random page reads. The fix is almost never "disable nested loops"; it is to fix the estimate (ANALYZE, extended statistics, rewriting the predicate) so the planner can see that a hash join is cheaper. The misestimates page walks through that diagnosis.
The reverse failure exists too: a hash join planned for a "huge" input that is actually tiny wastes a full scan of the build side, and a hash join planned for a build side that fits in memory may spill into dozens of batches when the estimate was 10x low.
Engine contrasts
- MySQL 8.0.18+ uses hash joins for equi-joins without a usable index (they replaced Block Nested Loop); it has no merge join. Memory is bounded by
join_buffer_size, beyond which it spills to disk files. Index nested loop ("ref"/"eq_ref" access in EXPLAIN) remains the default with an index. - SQLite uses nested loops only, but will build an automatic index on the inner table for the duration of the query when that is cheaper than repeated full scans (visible as
AUTOMATIC COVERING INDEXinEXPLAIN QUERY PLAN). - DuckDB uses parallel, vectorized hash joins for almost all equi-joins, with out-of-core partitioning when the build side exceeds memory; it also has range/inequality join operators (IEJoin, piecewise merge join) for band joins.
- PostgreSQL has all three plus Memoize and parallel hash join (workers share one hash table, PG 11+). Hash joins on outer joins can build on either side depending on the join type (
Hash Right Join,Hash Right Anti Joinin PG 16+).
Pitfalls
- Joining on mismatched types (
int=bigintis fine;text=intvia cast is not): the cast can prevent index use on the inner side, turning an index nested loop into a hash join over the full table. - Function-wrapped join keys (
lower(a.email) = lower(b.email)) need an expression index for index nested loops, and still work for hash joins. - Setting
work_memglobally high to avoid batches: every hash and sort node in every concurrent query may use that much (more for hashes). See the sizing math on the sorting page. - Reading
rowson the inner side of a nested loop as a total: it is per loop. - Disabling join types in production configs (
enable_nestloop = offglobally): it fixes one query and breaks hundreds of point lookups. - Forgetting duplicates: a many-to-many join can return far more rows than either input; estimate the output size, not just the inputs.
How the code was checked
- The work-count model ran under Python 3.13 and the cost formulas 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, including the
Batches: 8spill. The plans are real output; only the setup DDL echoes were omitted. - Page counts, buffer sizes and row widths in the models are example values.
Interview Q&A
When does a nested loop beat a hash join?
Answer
When the outer input is small and the inner side has an index on the join key: cost is proportional to the outer rows, it needs no memory, and it returns first rows immediately. Also when the condition is not an equality, since hash join cannot run it at all.
Explain how a hash join handles a build side that does not fit in memory.
Answer
It partitions both inputs by a hash of the join key into batches so that each build partition fits. One batch stays in memory; the others are written to temp files and joined batch by batch later. Each spilled row is written and read once more, about 3x the I/O of an in-memory hash join. In PostgreSQL EXPLAIN ANALYZE shows Batches: N on the Hash node.
Why would the planner choose a merge join?
Answer
When both inputs already arrive sorted on the join key (indexes or earlier operations), merge join is one pass with minimal memory and produces sorted output that a later ORDER BY or GROUP BY can use. It is also a robust choice for big-big joins when memory is tight.
A query that used to take 50 ms now takes 20 minutes, and the plan shows a nested loop with loops=2,000,000 on the inner side. What happened?
Answer
The planner estimated the outer side at a few rows, so a nested loop looked cheap; the real outer cardinality is two million. Compare estimated vs actual rows on the outer side, find why the estimate is low (stale stats, correlated predicates, new out-of-range values, generic plan), and fix the estimate so the planner picks a hash join.
What is Memoize in PostgreSQL plans?
Answer
A cache between a nested loop and its inner side keyed on the parameter (the outer join key). Repeated keys hit the cache instead of re-running the inner scan. It helps when the outer side has many duplicate keys and the inner lookup is expensive.
Which side of a hash join should be the build side?
Answer
The smaller one (by estimated bytes), because memory and batch count scale with the build side while the probe side is streamed. A wrong size estimate can put the large input on the build side.
Why does merge join appear when hash join is disabled in the real run?
Answer
Both inputs were available sorted on the join key through Index Only Scans on events_user_idx and users_pkey, so a merge join needed no Sort nodes.
How do MySQL and SQLite differ from PostgreSQL on joins?
Answer
MySQL 8.0.18+ uses hash joins for equi-joins without a usable index and has no merge join. SQLite uses nested loops only but may build an automatic transient index on the inner table.
What does a band join need?
Answer
A condition like b.id BETWEEN a.id AND a.id + 2 has no hash key, so it runs as a nested loop with an index range scan on the inner side, as the real plan shows.
Check yourself
Pick a two-table join from your schema. Run EXPLAIN with the planner's free choice, then with enable_hashjoin and enable_mergejoin off in turn (session only). Compare costs and note which index each alternative needed.
Elsewhere in the library
These pages stay as they are. This lesson only points at them: Indexes, Cardinality & EXPLAIN Plans, B-tree vs Hash vs GIN vs GiST, Composite & Covering Indexes, B-Tree Internals — Pages, Splits & Buffer Pool.
Go Deeper
- Use The Index, Luke: The Join Operation (nested loops, hash join, sort-merge)
- Use The Index, Luke: Nested Loops and the ORM N+1 problem
- Use The Index, Luke: Hash Join
- Use The Index, Luke: Sort-Merge Join
- PostgreSQL docs: Planner/Optimizer (join path generation)
- PostgreSQL docs: Resource consumption (work_mem, hash_mem_multiplier)
- MySQL 8.4 Reference Manual: Hash Join Optimization
- CMU 15-445/645 schedule (Joins Algorithms lecture and notes)
- CMU Database Group on YouTube