Sorting, Aggregation, Spills & Parallel Query - External Merge, Top-N, HashAggregate vs GroupAggregate & work_mem Math
Sort strategies (quicksort, external merge, top-N heapsort, incremental sort) with a runnable external merge sort pass counter and top-N vs full sort comparisons; HashAggregate vs GroupAggregate; real PG 17 spills (external merge Disk, HashAggregate Batches/Disk Usage, hash_mem_multiplier) and a Gather/Partial Aggregate parallel plan; runnable work_mem worst-case sizing math; temp-file monitoring; MySQL/SQLite/DuckDB/warehouse contrasts.
- 1Gist
- 2Maps
- 3Q&A
- 4Sandbox
Voice readout needs Web Speech Synthesis in this browser.
Question ladder
L1
What does work_mem limit?
Answer
The memory of each sort or hash operation before it spills, not the query or the connection.
L2
What does external merge Disk: 12256kB mean?
Answer
The sort did not fit in work_mem, so it wrote sorted runs to disk and merged them.
L3
Why is ORDER BY ... LIMIT 10 cheap without an index?
Answer
A bounded top-N heap keeps 10 rows, so memory is tiny and there is about one comparison per input row.
L4
HashAggregate or GroupAggregate?
Answer
Hash for unsorted input and moderate group counts; Group when the input is already sorted or groups are too many to hash.
L5
Why did a 1MB work_mem HashAggregate use 2105kB?
Answer
Hash nodes may use work_mem x hash_mem_multiplier, which defaults to 2.0.
L6
What does loops=3 on a Parallel Seq Scan mean?
Answer
Three processes (leader plus two workers) each scanned part of the table; rows are per process.
L7
How do you find spilling queries?
Answer
log_temp_files, pg_stat_database temp_files and temp_bytes, pg_stat_statements temp_blks_written, and EXPLAIN ANALYZE.
Failure modes
Hot query sorts on disk
An ORDER BY without a supporting index spills to an external merge on every call.
OOM from a high global work_mem
Concurrent reports with many nodes and parallel workers allocate far more than physical RAM, and the OOM killer ends backends.
Fewer workers than planned
A busy server launches fewer parallel workers, and the plan runs slower than expected.
Misconceptions
If the data is 12 MB on disk, 16 MB of work_mem is enough.
The in-memory sort used 27914kB because tuples carry pointers and per-tuple overhead.
Memory Usage in EXPLAIN is the query's total.
It is per node, and per process for parallel nodes.
Parallelism fixes a missing index.
Two workers make a full scan about 2-3x faster; an index makes a selective lookup thousands of times faster.
Interviewer traps
Raising work_mem globally because sorts spill.
Raise it per role or with SET LOCAL for the heavy statements, after measuring temp files.
Assuming LIMIT stops a HashAggregate early.
All input is aggregated before the LIMIT applies.
Design scenario
Same prompt for every reader.
Requirements
No OOM kills, API latency unaffected by reports, and reports that finish without excessive temp files.
Failure assumptions
- Several reports start at once.
- A report's group count is underestimated.
- Autoscaling adds API connections during peaks.
Constraints
- 16 GB shared_buffers.
- Reports must run on the primary for now.
Prompt
A 64 GB PostgreSQL server serves an OLTP API with 200 active queries and a reporting role that runs heavy GROUP BY queries. Size memory settings and parallelism.
API
Which roles and transactions get a larger work_mem, and how is it set?
Data
Which metrics and logs show spills and memory per query?
Architecture
How do max_parallel_workers_per_gather, connection limits and work_mem combine into a worst-case memory budget?
Overview
Sorts and aggregates are where a query's memory behaviour lives. PostgreSQL gives each sort or hash operation a budget of work_mem (hash-based operations get work_mem x hash_mem_multiplier, default 2.0); when the data does not fit, the operation spills to temp files and gets several times slower. Knowing the strategies lets you read a plan and predict its cost: quicksort in memory, external merge sort on disk, top-N heapsort when a LIMIT bounds the output, incremental sort when the input is already partly sorted; HashAggregate (one hash table entry per group, no sorted input needed) vs GroupAggregate (needs sorted input, streams with constant memory). On top sits parallel query: workers each scan a part of the table and run partial aggregates, a Gather (or Gather Merge) node collects the results, and every worker gets its own work_mem. The senior skill is the arithmetic: how much memory a query can really use, and when a spill is cheaper than the memory it would take to avoid it.
Aggregation strategies
| Strategy | Needs | Memory | Output order | Best when | Failure mode |
|---|---|---|---|---|---|
| HashAggregate | Hashable group keys | One entry per group (x hash_mem_multiplier budget) | Unordered | Few to moderate groups, unsorted input | Many groups: spills partitions to disk (PG 13+); pre-13 could exceed memory |
| GroupAggregate (sorted) | Input sorted on group keys | Constant (one group at a time) | Sorted by group keys | Input already sorted (index, merge join), or huge group counts | Needs a Sort below if not presorted; that sort may spill |
| Mixed / grouping sets | GROUPING SETS, ROLLUP, CUBE | Combination | Mixed | Multi-level summaries in one pass | Complex plans, multiple sorts |
| Partial + Finalize (parallel or partitionwise) | Aggregates that can combine partial states (count, sum, avg, min, max) | Per worker | Depends on final node | Large scans with parallel workers | Not every aggregate is parallel-safe (e.g. some user-defined, ordered-set aggregates) |
PostgreSQL 13 made HashAggregate spill-capable: when the hash table outgrows its budget, new groups are written to partitions on disk and processed in later batches, reported as Batches: N ... Disk Usage: .... Before 13, an underestimated group count could make a HashAggregate use far more memory than work_mem, a classic out-of-memory cause.
Raise work_mem globally or only where it is needed?
Prefer
Modest default, raised per role or transaction
Keep the global budget safe and give heavy jobs more with ALTER ROLE or SET LOCAL.
- The OLTP fleet's safe maximum was 27 MB for a 16 GB budget.
- Reporting could safely use 49 MB, or more for a single SET LOCAL job.
- Spills are found with log_temp_files and fixed one statement at a time.
Alternative
A large global work_mem
Set one big value so nothing spills.
- 256MB gave a 150 GB worst case for the OLTP fleet.
- One burst of concurrent reports can trigger the OOM killer.
- It hides missing indexes that should provide sort order.
How PostgreSQL picks a sort strategy
Diagram 1 condensed.
- 1
Already sorted?
An index or prior step can remove the Sort node entirely. - 2
Sorted on a prefix
Incremental Sort sorts small groups only. - 3
LIMIT above
Top-N heapsort keeps n rows in a bounded heap. - 4
Fits in work_mem
In-memory quicksort. - 5
Too big
External merge writes sorted runs to temp files and merges them.
Sorting strategies
Decisions
- 1
Step 1: ORDER BY / merge join / GROUP BY needs sorted input
- nextStep 2: input already sorted by an index or prior step?
- ?
Step 2: input already sorted by an index or prior step?
- fullyStep 3a: no Sort node at all
- on a key prefixStep 3b: Incremental Sort - sort small groups only
- noStep 3c: is there a LIMIT n above?
- 3
Step 3a: no Sort node at all
- 4
Step 3b: Incremental Sort - sort small groups only
- ?
Step 3c: is there a LIMIT n above?
- yes, n smallStep 4a: top-N heapsort - keep n rows in a bounded heap
- noStep 4b: fits in work_mem?
- 6
Step 4a: top-N heapsort - keep n rows in a bounded heap
- ?
Step 4b: fits in work_mem?
- yesStep 5a: in-memory quicksort
- noStep 5b: external merge - write sorted runs to temp files, k-way merge
- 8
Step 5a: in-memory quicksort
- 9
Step 5b: external merge - write sorted runs to temp files, k-way merge
Lesson map
Sorting, Aggregation, Spills & Parallel Query - External Merge, Top-N, HashAggregate vs GroupAggregate & work_mem Math
Sort strategies (quicksort, external merge, top-N heapsort, incremental sort) with a runnable external merge sort pass counter and top-N vs full sort comparisons; HashAggregate vs GroupAggregate; real PG 17 spills (external merge Disk, HashAggregate Batches/Disk Usage, hash_mem_multiplier) and a Gather/Partial Aggregate parallel plan; runnable work_mem worst-case sizing math; temp-file monitoring; MySQL/SQLite/DuckDB/warehouse 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: ORDER BY / merge join / GROUP BY needs sorted input"] b["Step 2: input already sorted by an index or prior step?"] n["Step 3a: no Sort node at all"] i["Step 3b: Incremental Sort - sort small groups only"] l["Step 3c: is there a LIMIT n above?"] t["Step 4a: top-N heapsort - keep n rows in a bounded heap"] m["Step 4b: fits in work_mem?"] q["Step 5a: in-memory quicksort"] x["Step 5b: external merge - write sorted runs to temp files, k-way merge"] a -->|continues| b b -->|fully| n b -->|on a key prefix| i b -->|no| l l -->|yes, n small| t l -->|no| m m -->|yes| q m -->|no| x
External merge sort is the classic out-of-core algorithm: read as much as fits in memory, sort it, write it out as a run; repeat; then merge runs with a k-way heap merge. With memory for k-way merges, the number of merge passes is about ceil(log_k(number of runs)); each pass reads and writes all the data once. The runnable model counts runs, passes and spilled rows for different memory budgets and fan-ins, then compares a bounded top-N heap against a full sort.
"""External merge sort vs top-N heapsort, the two strategies behind ORDER BY.
- External merge sort: sort memory-sized runs, spill them, then k-way merge (what PostgreSQL reports
as "Sort Method: external merge Disk: ...").
- Top-N heapsort: ORDER BY ... LIMIT n keeps only n rows in a bounded heap ("top-N heapsort").
Memory budgets and sizes are example values (rows, not bytes, to keep it simple).
"""
import heapq, math, random
random.seed(7)
data = [random.randrange(1_000_000) for _ in range(200_000)]
compares = 0
class Key: # wrapper that counts comparisons
__slots__ = ("v",)
def __init__(self, v): self.v = v
def __lt__(self, other):
global compares; compares += 1; return self.v < other.v
def external_sort(xs, mem_rows, fan_in):
# Phase 1: run generation. Each run fits in "work_mem" and is written to a temp "file" (a list here).
runs = [sorted(xs[i:i + mem_rows]) for i in range(0, len(xs), mem_rows)]
initial, passes, spilled = len(runs), 0, len(xs) # every row written once in phase 1
# Phase 2: merge at most fan_in runs at a time until one run remains.
while len(runs) > 1:
runs = [list(heapq.merge(*runs[i:i + fan_in])) for i in range(0, len(runs), fan_in)]
passes += 1; spilled += len(xs) if len(runs) > 1 else 0 # intermediate passes re-spill everything
return runs[0], initial, passes, spilled
for mem_rows, fan_in in [(250_000, 8), (20_000, 8), (2_000, 8), (2_000, 64)]:
out, initial, passes, spilled = external_sort(data, mem_rows, fan_in)
assert out == sorted(data)
predicted = 0 if initial == 1 else math.ceil(math.log(initial, fan_in))
print(f"mem={mem_rows:>7,} rows fan-in={fan_in:>2}: runs={initial:>3} merge passes={passes} "
f"(predicted {predicted}) rows written to temp={spilled if initial > 1 else 0:,}")
# Top-N: ORDER BY v DESC LIMIT 10 without sorting everything
compares = 0
top = heapq.nlargest(10, (Key(v) for v in data)) # bounded heap of size 10
top_n_compares = compares
compares = 0
full = sorted((Key(v) for v in data), reverse=True)[:10] # full sort, then take 10
print("top-10 identical:", [k.v for k in top] == [k.v for k in full])
print(f"comparisons: top-N heap={top_n_compares:,} vs full sort={compares:,} "
f"(~{compares / top_n_compares:.0f}x more); top-N memory = 10 rows, full sort = {len(data):,} rows")Output:
mem=250,000 rows fan-in= 8: runs= 1 merge passes=0 (predicted 0) rows written to temp=0
mem= 20,000 rows fan-in= 8: runs= 10 merge passes=2 (predicted 2) rows written to temp=400,000
mem= 2,000 rows fan-in= 8: runs=100 merge passes=3 (predicted 3) rows written to temp=600,000
mem= 2,000 rows fan-in=64: runs=100 merge passes=2 (predicted 2) rows written to temp=400,000
top-10 identical: True
comparisons: top-N heap=200,413 vs full sort=3,263,634 (~16x more); top-N memory = 10 rows, full sort = 200,000 rowsTwo practical consequences. First, going from "fits in memory" to "slightly too big" is a cliff (all data written and read at least once), but going from 10 runs to 100 runs is a gentle slope (one more pass): a sort that spills a little is often fine, and doubling work_mem to avoid it may not be worth the memory risk. Second, top-N is not a sort: ORDER BY ... LIMIT 10 keeps ten rows, needs almost no memory and does about one comparison per input row, which is why PostgreSQL reports top-N heapsort Memory: 25kB even under a tiny work_mem. (With an index on the sort key, there is no sort at all: the index scan delivers rows in order and the LIMIT stops it early.)
Real PostgreSQL: sort methods, aggregation strategies, spills and Gather
-- Page 5: sort methods (quicksort / external merge / top-N heapsort), HashAggregate vs GroupAggregate,
-- HashAggregate spilling, and a parallel plan with Gather.
\pset footer off
SET client_min_messages = warning;
DROP SCHEMA IF EXISTS p5 CASCADE;
CREATE SCHEMA p5;
SET search_path = p5;
SET max_parallel_workers_per_gather = 0;
CREATE TABLE sales AS
SELECT g AS id, ((g::bigint * 7919) % 100000)::int AS customer_id, (g * 37) % 1000 + 0.5 AS amount
FROM generate_series(1, 500000) g;
VACUUM ANALYZE sales;
-- 1) Full sort that does not fit in work_mem -> external merge on disk
SET work_mem = '1MB';
EXPLAIN (ANALYZE, COSTS OFF, TIMING OFF, SUMMARY OFF) SELECT * FROM sales ORDER BY amount, id;
-- 2) Same sort with enough memory -> in-memory quicksort
SET work_mem = '64MB';
EXPLAIN (ANALYZE, COSTS OFF, TIMING OFF, SUMMARY OFF) SELECT * FROM sales ORDER BY amount, id;
-- 3) ORDER BY ... LIMIT -> bounded top-N heapsort, tiny memory even with small work_mem
SET work_mem = '1MB';
EXPLAIN (ANALYZE, COSTS OFF, TIMING OFF, SUMMARY OFF) SELECT * FROM sales ORDER BY amount DESC, id LIMIT 10;
-- 4) Aggregation strategy: few groups -> HashAggregate; force sort-based GroupAggregate for comparison
EXPLAIN (COSTS OFF) SELECT amount, count(*) FROM sales GROUP BY amount;
SET enable_hashagg = off;
EXPLAIN (COSTS OFF) SELECT amount, count(*) FROM sales GROUP BY amount;
RESET enable_hashagg;
-- 5) Many groups + small work_mem -> HashAggregate spills partitions to disk (PostgreSQL 13+)
SET work_mem = '1MB';
EXPLAIN (ANALYZE, COSTS OFF, TIMING OFF, SUMMARY OFF) SELECT customer_id, sum(amount) FROM sales GROUP BY customer_id;
-- 6) Parallel query: allow workers and make them cheap for the demo -> Gather + Partial/Finalize aggregate
SET max_parallel_workers_per_gather = 2;
SET parallel_setup_cost = 0; SET min_parallel_table_scan_size = 0;
SET work_mem = '64MB';
EXPLAIN (ANALYZE, COSTS OFF, TIMING OFF, SUMMARY OFF) SELECT count(*), sum(amount) FROM sales WHERE amount > 100;Real output (PostgreSQL 17.11, local sandbox, psql; setup DDL echoes omitted):
SET max_parallel_workers_per_gather = 0;
SET work_mem = '1MB';
EXPLAIN (ANALYZE, COSTS OFF, TIMING OFF, SUMMARY OFF) SELECT * FROM sales ORDER BY amount, id;
QUERY PLAN
------------------------------------------------------
Sort (actual rows=500000 loops=1)
Sort Key: amount, id
Sort Method: external merge Disk: 12256kB
-> Seq Scan on sales (actual rows=500000 loops=1)
SET work_mem = '64MB';
EXPLAIN (ANALYZE, COSTS OFF, TIMING OFF, SUMMARY OFF) SELECT * FROM sales ORDER BY amount, id;
QUERY PLAN
------------------------------------------------------
Sort (actual rows=500000 loops=1)
Sort Key: amount, id
Sort Method: quicksort Memory: 27914kB
-> Seq Scan on sales (actual rows=500000 loops=1)
SET work_mem = '1MB';
EXPLAIN (ANALYZE, COSTS OFF, TIMING OFF, SUMMARY OFF) SELECT * FROM sales ORDER BY amount DESC, id LIMIT 10;
QUERY PLAN
------------------------------------------------------------
Limit (actual rows=10 loops=1)
-> Sort (actual rows=10 loops=1)
Sort Key: amount DESC, id
Sort Method: top-N heapsort Memory: 25kB
-> Seq Scan on sales (actual rows=500000 loops=1)
EXPLAIN (COSTS OFF) SELECT amount, count(*) FROM sales GROUP BY amount;
QUERY PLAN
-------------------------
HashAggregate
Group Key: amount
-> Seq Scan on sales
SET enable_hashagg = off;
EXPLAIN (COSTS OFF) SELECT amount, count(*) FROM sales GROUP BY amount;
QUERY PLAN
-------------------------------
GroupAggregate
Group Key: amount
-> Sort
Sort Key: amount
-> Seq Scan on sales
RESET enable_hashagg;
SET work_mem = '1MB';
EXPLAIN (ANALYZE, COSTS OFF, TIMING OFF, SUMMARY OFF) SELECT customer_id, sum(amount) FROM sales GROUP BY customer_id;
QUERY PLAN
----------------------------------------------------------
HashAggregate (actual rows=100000 loops=1)
Group Key: customer_id
Batches: 81 Memory Usage: 2105kB Disk Usage: 15288kB
-> Seq Scan on sales (actual rows=500000 loops=1)
SET max_parallel_workers_per_gather = 2;
SET parallel_setup_cost = 0;
SET min_parallel_table_scan_size = 0;
SET work_mem = '64MB';
EXPLAIN (ANALYZE, COSTS OFF, TIMING OFF, SUMMARY OFF) SELECT count(*), sum(amount) FROM sales WHERE amount > 100;
QUERY PLAN
---------------------------------------------------------------------------
Finalize Aggregate (actual rows=1 loops=1)
-> Gather (actual rows=3 loops=1)
Workers Planned: 2
Workers Launched: 2
-> Partial Aggregate (actual rows=1 loops=3)
-> Parallel Seq Scan on sales (actual rows=150000 loops=3)
Filter: (amount > '100'::numeric)
Rows Removed by Filter: 16667What each section shows:
- external merge Disk: 12256kB with
work_mem = 1MB: 500,000 rows did not fit, so the sort produced runs on disk and merged them. - quicksort Memory: 27914kB with
work_mem = 64MB: same sort, entirely in memory. Notice that the in-memory footprint (about 27 MB) is larger than the on-disk footprint (about 12 MB): in-memory tuples carry pointers and per-tuple overhead, so "the data is 12 MB on disk, 16 MB of work_mem should do" is wrong. - top-N heapsort Memory: 25kB with
LIMIT 10: reads all 500,000 rows but keeps only 10. - HashAggregate vs GroupAggregate for 1,000 groups: the planner prefers hashing; disabling it yields Sort + GroupAggregate. With an index on
amount, GroupAggregate over an Index Scan would need no sort at all. - HashAggregate spill: 100,000 groups with
work_mem = 1MBgiveBatches: 81 Memory Usage: 2105kB Disk Usage: 15288kB. Memory used is about 2 MB, not 1 MB, because hash nodes may usework_mem x hash_mem_multiplier(2.0). - Parallel plan:
Finalize Aggregate <- Gather (Workers Launched: 2) <- Partial Aggregate <- Parallel Seq Scan (loops=3). Three processes (leader plus two workers) each scanned about a third of the table (actual rows=150000 loops=3means per process) and produced one partial result each; the leader combined three rows. The demo loweredparallel_setup_costandmin_parallel_table_scan_sizeso a small table qualifies; on real tables the defaults decide.
work_mem sizing: the arithmetic interviewers want
work_mem is per operation, not per query or per connection. A query with three sorts and two hash joins can use five budgets; with two parallel workers, each worker gets its own budgets too; hash operations get twice the budget by default. The worst case is roughly:
concurrent active queries x memory-hungry nodes per query x (1 + workers) x work_mem (x hash_mem_multiplier for hash nodes)
// work_mem is a per-operation budget, not per-query or per-connection.
// Worst case memory ~= active queries * memory-hungry nodes per query * (1 + parallel workers) * budget,
// where hash-based nodes (Hash, HashAggregate) may use work_mem * hash_mem_multiplier (default 2.0, PG 15+).
// Every number below is an example value for a sizing exercise.
interface Workload {
name: string;
activeQueries: number; // concurrently executing, not just connected
sortNodes: number; // Sort / Incremental Sort / Materialize per query
hashNodes: number; // Hash (join build) / HashAggregate per query
parallelWorkers: number; // workers per Gather; each gets its own budget per node
}
const GB = 1024 * 1024; // kB per GB
function worstCaseGb(w: Workload, workMemKb: number, hashMult = 2.0): number {
const perProcess = w.sortNodes * workMemKb + w.hashNodes * workMemKb * hashMult;
return (w.activeQueries * (1 + w.parallelWorkers) * perProcess) / GB;
}
// Largest work_mem that keeps the worst case inside a memory budget left after shared_buffers & OS cache.
function maxWorkMemKb(w: Workload, budgetGb: number, hashMult = 2.0): number {
const unitsPerKb = w.activeQueries * (1 + w.parallelWorkers) * (w.sortNodes + w.hashNodes * hashMult);
return Math.floor((budgetGb * GB) / unitsPerKb);
}
const oltp: Workload = { name: "OLTP API", activeQueries: 200, sortNodes: 1, hashNodes: 1, parallelWorkers: 0 };
const report: Workload = { name: "reporting", activeQueries: 10, sortNodes: 3, hashNodes: 4, parallelWorkers: 2 };
const budgetGb = 16; // e.g. 64 GB box: 16 GB shared_buffers, ~32 GB OS cache, 16 GB for query memory
for (const w of [oltp, report]) {
for (const wm of [4 * 1024, 64 * 1024, 256 * 1024]) {
console.log(`${w.name.padEnd(9)} work_mem=${String(wm / 1024).padStart(3)}MB -> worst case ${worstCaseGb(w, wm).toFixed(1)} GB`);
}
console.log(`${w.name.padEnd(9)} max safe work_mem for ${budgetGb} GB: ${Math.floor(maxWorkMemKb(w, budgetGb) / 1024)} MB`);
}
console.log("pattern: keep the global default modest, raise it per role/session/transaction for known heavy jobs:");
console.log(" SET LOCAL work_mem = '256MB'; -- inside the reporting transaction only");Output:
OLTP API work_mem= 4MB -> worst case 2.3 GB
OLTP API work_mem= 64MB -> worst case 37.5 GB
OLTP API work_mem=256MB -> worst case 150.0 GB
OLTP API max safe work_mem for 16 GB: 27 MB
reporting work_mem= 4MB -> worst case 1.3 GB
reporting work_mem= 64MB -> worst case 20.6 GB
reporting work_mem=256MB -> worst case 82.5 GB
reporting max safe work_mem for 16 GB: 49 MB
pattern: keep the global default modest, raise it per role/session/transaction for known heavy jobs:
SET LOCAL work_mem = '256MB'; -- inside the reporting transaction onlyExpectedOLTP API work_mem= 4MB -> worst case 2.3 GB OLTP API work_mem= 64MB -> worst case 37.5 GB OLTP API work_mem=256MB -> worst case 150.0 GB OLTP API max safe work_mem for 16 GB: 27 MB reporting work_mem= 4MB -> worst case 1.3 GB reporting work_mem= 64MB -> worst case 20.6 GB reporting work_mem=256MB -> worst case 82.5 GB reporting max safe work_mem for 16 GB: 49 MB pattern: keep the global default modest, raise it per role/session/transaction for known heavy jobs: SET LOCAL work_mem = '256MB'; -- inside the reporting transaction only
Press Run. Snippets must be self-contained — no network, files, or native modules.
The numbers are example values for a hypothetical 64 GB server, but the shape is universal: an OLTP fleet with 200 active queries cannot afford a large global work_mem even though each query is simple, while a handful of reporting sessions can use much more, especially if raised only where needed (ALTER ROLE reporting SET work_mem = '256MB', or SET LOCAL work_mem inside one transaction). Real usage is usually far below the worst case (not every node uses its full budget at the same time), which is why teams set it by measurement: watch temp_files and temp_bytes in pg_stat_database, set log_temp_files = 0 (or a size threshold) to log every spill with its query, and raise budgets only for the statements that spill on a hot path.
| Symptom | Likely cause | Fix |
|---|---|---|
Sort Method: external merge on a hot query | Sort bigger than work_mem | Index providing the order; LIMIT/top-N; raise work_mem for that role or transaction |
Batches: N > 1 on Hash or HashAggregate | Build side or group count bigger than budget | Fix row estimate (it may be planned as smaller), raise budget locally, pre-aggregate |
| OOM kills during reports | Too-high global work_mem x concurrency x parallelism | Lower global value, raise per role; cap max_parallel_workers_per_gather |
temp_bytes grows steadily | Many queries spilling a little | Find them via log_temp_files; often a missing index for ORDER BY |
| Parallel plan slower than serial | Startup overhead, Gather bottleneck, skewed partitions | Leave parallel to big scans; check Workers Launched vs Planned |
Parallel query in one picture
Flow
- 1
Step 5: Finalize Aggregate in the leader - combine partial states
- nextStep 4: Gather - collect rows from workers, unordered (Gather Merge keeps order)
- 2
Step 4: Gather - collect rows from workers, unordered (Gather Merge keeps order)
- nextStep 3a: Worker 1 - Partial Aggregate
- nextStep 3b: Worker 2 - Partial Aggregate
- nextStep 3c: Leader also participates - Partial Aggregate
- 3
Step 3a: Worker 1 - Partial Aggregate
- nextStep 2a: Parallel Seq Scan - grabs blocks from a shared counter
- 4
Step 3b: Worker 2 - Partial Aggregate
- nextStep 2b: Parallel Seq Scan
- 5
Step 3c: Leader also participates - Partial Aggregate
- nextStep 2c: Parallel Seq Scan
- 6
Step 2a: Parallel Seq Scan - grabs blocks from a shared counter
- relatedStep 1: table blocks handed out dynamically
- 7
Step 2b: Parallel Seq Scan
- relatedStep 1: table blocks handed out dynamically
- 8
Step 2c: Parallel Seq Scan
- relatedStep 1: table blocks handed out dynamically
- 9
Step 1: table blocks handed out dynamically
Parallel-aware nodes in PostgreSQL include Parallel Seq Scan, Parallel Index (Only) Scan, Parallel Bitmap Heap Scan, Parallel Hash Join (one shared hash table), Parallel Append and partial aggregates. Limits: max_parallel_workers_per_gather (default 2), max_parallel_workers (default 8) and max_worker_processes (default 8), so a busy server may launch fewer workers than planned; queries that write data (apart from some CREATE TABLE AS and CREATE INDEX paths), use parallel-unsafe functions, or run inside cursors generally run serially.
Engine contrasts
- MySQL 8: sorts use
sort_buffer_sizeper sort and spill to temp files ("Using filesort" in EXPLAIN); internal temporary tables for GROUP BY use the TempTable engine withtmp_table_size/temptable_max_ram, then spill to disk. Parallelism is limited (parallel clustered-index reads for some operations likeCHECK TABLEandSELECT COUNT(*)), not a general parallel executor. - SQLite sorts with an external merge sorter bounded by its cache and temp store settings; it is single-threaded per query (with optional worker threads for sorting via
PRAGMA threads). - DuckDB is parallel by default (morsel-driven), uses vectorized hash aggregates with out-of-core spilling, and caps total memory with one global
memory_limitinstead of per-operator budgets, which avoids the "work_mem x everything" arithmetic. - Warehouses (BigQuery, Snowflake, Redshift) spill to local SSD and then remote storage; query profiles show "bytes spilled" per stage, the same concept at a larger scale.
What happens if you choose otherwise
- Global
work_mem = 1GB"because sorts were spilling": one burst of concurrent reports can allocate far more than physical RAM and the kernel OOM killer takes out the postmaster's children (and with them every connection). - Tiny
work_memon an analytics box: every sort and hash spills, queries become I/O bound on temp files, and the planner may choose worse strategies because it prices the spills. - Relying on parallelism to fix a missing index: two workers make a full scan about 2-3x faster; an index makes a selective lookup thousands of times faster.
- GROUP BY on a high-cardinality key with a big LIMIT expectation: without an index providing order, the whole input is still aggregated before the LIMIT applies.
Pitfalls
- Reading
Memory Usagein EXPLAIN as the query's total memory: it is per node (and per process for parallel nodes). - Forgetting
hash_mem_multiplierwhen interpreting spills: hash nodes can legitimately use 2xwork_mem. - Assuming
LIMITalways makes a query cheap: above a HashAggregate or full Sort, all input is still processed. - Measuring sort speed on a warm cache, then running on a cold cache in production.
- Expecting
Workers Planned=Workers Launched: on a busy server, fewer workers may be available. DISTINCTandUNION(without ALL) also sort or hash: they cost like a GROUP BY.
How the code was checked
- The external sort model ran under Python 3.13 and the work_mem sizing 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 external merge, HashAggregate spill and Gather plan. Only the setup DDL echoes were omitted.
- Memory budgets in both models are example values for a sizing exercise.
Interview Q&A
What does work_mem control and why can a single query use many multiples of it?
Answer
It is the memory budget for each sort or hash operation before spilling to disk. A query can contain several such nodes, each parallel worker gets its own, and hash nodes may use work_mem x hash_mem_multiplier. So memory scales with nodes, workers and concurrent queries.
How does an external merge sort work and what does it cost?
Answer
Sort memory-sized chunks into runs written to disk, then k-way merge the runs (possibly in multiple passes). Each pass reads and writes all the data, so the cost is about 2 x data size x number of passes; passes grow logarithmically with the number of runs.
Why is ORDER BY ... LIMIT 10 cheap even without an index?
Answer
The executor uses a bounded top-N heap: it scans everything but keeps only 10 rows, with tiny memory and about one comparison per row. With an index on the sort key it does not even scan everything.
HashAggregate vs GroupAggregate?
Answer
HashAggregate builds a hash table of groups from unsorted input (fast, unordered output, memory grows with group count, spills in PG 13+). GroupAggregate needs input sorted on the group keys and streams one group at a time with constant memory; it wins when the order comes for free from an index or a merge join, or when groups are too many to hash.
Explain the parallel plan Finalize Aggregate <- Gather <- Partial Aggregate <- Parallel Seq Scan.
Answer
Workers and the leader each scan a dynamic share of the table's blocks and compute a partial aggregate state; Gather collects those partial states; the leader combines them in Finalize Aggregate. Row counts on parallel nodes are per process (loops = number of processes).
How do you find which queries spill?
Answer
log_temp_files (log every temp file above a size with its statement), pg_stat_database.temp_files/temp_bytes for trends, EXPLAIN ANALYZE for external merge, Disk Usage or Batches > 1, and pg_stat_statements' temp_blks_written per query.
Why is a small spill often acceptable?
Answer
Going from fits to slightly too big writes and reads the data once, but going from 10 runs to 100 runs adds only one merge pass, so doubling work_mem to avoid a small spill may not be worth the memory risk.
How does DuckDB avoid the work_mem arithmetic?
Answer
It caps total memory with one global memory_limit and spills out of core, instead of giving each operator its own budget.
Check yourself
Turn on log_temp_files = 0 on a staging database for an hour of realistic traffic. List the three statements that wrote the most temp data and decide for each whether an index, a LIMIT or a local work_mem increase is the right fix.
Elsewhere in the library
These pages stay as they are. This lesson only points at them: SQL Analytics — Window Functions, CTEs & Set-Based Thinking, Memory Profiling — Allocations, Leaks & GC Pressure, Performance Engineering — Profiling, Flame Graphs, Latency Budgets & Hot Paths, EXPLAIN & EXPLAIN ANALYZE.
Go Deeper
- PostgreSQL docs: Resource Consumption (work_mem, hash_mem_multiplier, temp files)
- PostgreSQL docs: Parallel Query
- PostgreSQL docs: Using EXPLAIN (sort and hash details)
- PostgreSQL docs: Error reporting and logging (log_temp_files)
- PostgreSQL docs: Cumulative statistics (pg_stat_database temp_files, temp_bytes)
- Use The Index, Luke: Top-N queries and pipelined ORDER BY
- DuckDB: EXPLAIN ANALYZE and operator profiling
- CMU 15-445/645 schedule (Sorting and Aggregation lecture)
- CMU Database Group on YouTube