What happens when an operator that must buffer its input — such as a sort or a hash-join build side — is granted less memory than its input needs, and how do you recognize and fix that in production?
answer
- sort → sorted runs + merge passes
- hash → partition to temp, join per partition
- estimate vs actual at the breaker
- width × rows drives grant; project fewer columns
- grant × workers × concurrency = real budget
basics
~20 sIt spills: the sort writes sorted runs to temporary storage and merges them later; the hash build partitions the input to disk and processes partitions in turn. The query stays correct but does extra passes of I/O. Diagnose via spill counters and temp-file usage; fix the cardinality estimate, the memory grant, or the plan.
solid answer
~60 sThe operator degrades gracefully rather than failing. A sort becomes an **external merge sort**: fill memory, sort, write a run, repeat, then merge runs — one extra read and write of the data per merge pass. A hash build becomes a **partitioned (grace) hash join**: both inputs are hashed into partitions, partitions that do not fit are written to temp files, and each pair is joined in a later pass; a partition that is still too large is re-partitioned recursively, and badly skewed keys can defeat that. Symptoms: temporary-file bytes or spill counters in the plan's runtime statistics, a large gap between estimated and actual rows at the breaker, sudden I/O on the temp volume, and runtimes that jump nonlinearly when data grows slightly. Fixes in order of leverage: correct the estimate (refresh statistics, fix a predicate the estimator cannot see, break a correlated multi-column assumption); reduce the volume (project fewer columns, filter earlier, top-N instead of full sort); reshape the plan (build on the smaller side, use an index-supplied order); and only then raise the working-memory grant — remembering that the grant multiplies by concurrency.
code
text · 6 linesSort (est rows=5,000 actual rows=4,812,004)
Sort Method: external merge Disk: 1,842,560 kB
Sort Key: o.created_at
-> Hash Join (est rows=5,000 actual rows=4,812,004)
Hash Batches: 16 Memory Usage: 4,096 kB
(batches > 1 => the build side spilled and was partitioned)go deeper
Know that when a sort or hash operator does not fit in memory it writes temporary files to disk and gets slower, and that the result is still correct.
Describe external merge sort and partitioned hash join, and name the obvious levers: fewer columns, better statistics, more working memory.
Diagnose from runtime plan statistics and estimate-versus-actual gaps, handle skew explicitly, and reason about grant multiplied by workers and concurrency before touching settings.
Decide when spilling is the correct designed behaviour to tune around versus a symptom of estimation failure, and set workload-level memory and concurrency policy accordingly.
## Why a spill happens at all Blocking operators must hold state proportional to their input: a sort holds the rows, a hash join's build side holds a hash table, a hash aggregate holds one entry per group. The engine gives the operator a **memory grant** — a per-operation working-memory budget derived from configuration and from the optimizer's estimated row count and row width. If the real data exceeds that budget, the engine must either fail the query or move part of the state to temporary storage. Every serious engine chooses the latter: correctness is preserved, performance degrades. ## What spilling looks like mechanically **External merge sort.** Read input until memory is full, sort that chunk in memory, write it out as a sorted *run*, and repeat. When input ends, merge the runs. With enough memory for a single merge pass, the total extra I/O is one write plus one read of the whole data set. With too little memory to merge all runs at once, you need multiple merge passes and the I/O multiplies. Cost grows roughly with the number of passes, i.e. logarithmically in the ratio of data size to memory — but each pass is a full sequential scan of a large temp file. **Partitioned (grace) hash join.** Hash the build input into P partitions by a partitioning hash of the join key. Partitions that fit stay in memory; the rest are written to temp files. The probe input is partitioned the same way. Then, for each spilled partition pair, load the build partition into memory and probe it. If one partition is still too large, recurse with a different hash function. This is why **skew is the pathological case**: if a single key value accounts for a huge share of rows, no partitioning function separates it, recursion does not shrink it, and the engine falls back to a much worse strategy or thrashes. **Hash aggregation** spills similarly by partitioning groups, and is likewise vulnerable to a few enormous groups. ## Why the query gets so much slower - Extra full read+write passes over data that was previously CPU-resident. - Temp I/O is often on shared storage, so several spilling queries contend for the same device. - Spill I/O is not free of CPU either: serialization, partitioning hashes, and buffer copying. - The transition is abrupt. Just under the grant, everything is in memory; just over, you pay whole extra passes. Runtimes therefore appear to jump discontinuously when a table grows slightly or a parameter changes selectivity. ## Recognizing it 1. **Runtime plan statistics.** Look for spill or temp-file indicators on the breaker: bytes written to temporary storage, number of sort runs or merge passes, number of hash batches or recursion level. 2. **Estimate versus actual rows** at the blocking operator. Spills almost always accompany a large underestimate: stale statistics, correlated predicates treated as independent, a filter on an expression the estimator cannot reason about, or a parameterized plan reused for an atypical parameter value. 3. **System-level signals.** Temp tablespace or temp directory growth, I/O on the temp volume, waits attributed to temp reads/writes. 4. **Workload shape.** Nonlinear slowdowns and high variance for the same query text across parameters. ## Fixing it, in order of leverage 1. **Fix the estimate.** Refresh statistics; add multi-column or expression statistics where the estimator's independence assumption is wrong; rewrite a predicate into a form the estimator can use. A correct estimate often produces both a bigger grant and a better plan. 2. **Reduce the volume that must be buffered.** Project only the needed columns (row width is a direct multiplier on sort memory), push filters below the breaker, aggregate earlier, and use a top-N sort when only a few rows are needed so memory becomes O(N). 3. **Reshape the plan.** Make the smaller input the build side. Provide an index that supplies the required ordering so the sort disappears entirely. Consider a merge join over already-ordered inputs, which streams instead of buffering. Pre-aggregate to shrink a grouping input. 4. **Address skew explicitly.** Separate the heavy key (handle it with its own path), or salt/split it, when a single value dominates a join or grouping key. 5. **Raise the memory grant — last, and deliberately.** Working memory is per operation, and a plan may contain several blocking operators, each under multiple parallel workers, across many concurrent sessions. The real budget is grant × operators × workers × concurrency. Raising it globally to fix one report is how machines run out of memory. Prefer raising it for the specific session or workload that needs it. ## What good judgement sounds like Spilling is not inherently a bug: for a genuinely large sort it is the correct, designed behaviour, and giving every query enough memory to avoid it is unaffordable. The judgement call is whether this operator *should* have been that big. If actual rows match the estimate and the data really is huge, tune the spill (fewer, larger runs; faster temp storage) and accept it. If the actual dwarfs the estimate, the estimate is the bug and the spill is only the symptom.
- Why is data skew especially damaging to a spilling hash join?Partitioned hash joins assume a hash function spreads keys evenly, so an oversized partition can be split further with a different hash. A single dominant key value hashes to one partition no matter what function is used, so recursive partitioning never shrinks it. The engine then either thrashes or falls back to a much slower strategy for that partition. The fix is to handle the heavy key separately rather than to add memory.
- Is raising the per-operation working-memory setting a good first response to spilling?Usually not. The setting applies per blocking operator, per parallel worker, per concurrent session, so a global increase multiplies into far more memory than the number suggests and risks exhausting the machine. Check first whether the operator is large because the estimate was wrong or because the plan buffers more than it needs; fix that, or raise the grant only for the specific session or workload that genuinely needs it.
- How can you tell a spill from ordinary slow I/O?Look at the operator's runtime statistics rather than at the system alone: sort method and disk bytes, hash batch or recursion counts, and temporary-file usage attributed to the query. Pair that with the estimated-versus-actual row counts at the same operator. Slow base-table I/O shows up at the scans with no temp activity, while a spill shows temp writes concentrated at the blocking operator.
Sorting a card index on a desk that is too small: you sort as many cards as fit, stack that pile on the floor, and repeat — then you have to merge all the floor piles afterwards, walking the whole room again.
saying these in an interview costs you the question
- Treating spilling as data corruption or an error rather than a designed, correct fallback.
- Reflexively raising the global working-memory setting without accounting for concurrency and parallel workers.
- Assuming a spilled hash join can always be fixed by more partitions, ignoring single-key skew.
- Ignoring row width — sorting unnecessary wide columns is a common and easily removed cause.
- Diagnosing from total runtime alone and never comparing estimated versus actual rows at the blocking operator.