skip to content

When a database sorts or groups a large result set, the plan may report that the operator "spilled to disk". What does spilling mean, what triggers it, and how does it show up in query performance?

level: juniorimportance: must knowfreq 60%

answer

  1. Blocking operator needs state → state needs RAM
  2. Budget is per operator, not per query
  3. Sort spill = sorted runs + merge; hash spill = partitions
  4. Hash agg memory ∝ distinct groups, not input rows
  5. Step change in latency, not a gradual slope

basics

~20 s

Each sort or grouping operator gets a limited memory budget. If the rows it must hold exceed that budget, the engine writes partial results into temporary files and reads them back later. Spilling replaces memory work with disk I/O, so latency jumps sharply.

solid answer

~60 s

Blocking operators like sort and hash aggregation need to hold state: a sort needs all input rows before it can emit the first one, and a hash aggregate needs one entry per distinct group. The engine gives each such operator a memory budget (PostgreSQL calls it `work_mem`, SQL Server calls it a memory grant, Oracle uses PGA workarea sizing). While the state fits, everything happens in RAM. When it does not, the operator switches strategy: a sort writes sorted runs to temporary files and merges them, a hash aggregate partitions the input and pushes some partitions to disk to reprocess later. The practical effect is a step change, not a gentle slope — the same query on slightly more data can go from 200 ms to 20 s. Symptoms are temp-file bytes in the plan, temp-space I/O in OS metrics, and elapsed time that grows far faster than row count. Fixes: reduce rows or width before the operator, use an index that supplies order, or raise the budget for that workload.

go deeper

for a junior

Be able to say what spilling is (operator state exceeded its memory budget, so it uses temp files) and that it makes the query much slower.

for a middle

Add the mechanics: which operators are stateful, that the budget is per operator, and that hash aggregation scales with distinct groups while sorting scales with rows carried.

for a senior

Show how you detect it in plans and system metrics, and give a remediation order — reduce rows/columns, supply order via an index, fix estimates, then tune memory for the specific workload.

for a principal

Frame it as a memory-admission problem: per-operation budgets times plan shape times concurrency is the real exposure, and the policy question is how much memory a single query may claim before the engine degrades gracefully instead of the machine falling over.

## What "spilling" actually means A query plan is a tree of operators. Most are *streaming*: a filter looks at one row, decides, and passes it on, holding nothing. Some are *stateful*: they must accumulate data before they can produce correct output. Sorting is the classic case — you cannot know which row is smallest until you have seen every row. Grouping by hash is another — you must keep one accumulator per distinct group key until the input ends. That accumulated data has to live somewhere. Every engine gives such an operator a bounded slice of RAM. If the state fits in the slice, the operator runs entirely in memory. If it does not, the operator cannot simply allocate more (that would let one query exhaust the machine), so it *spills*: it writes part of its state into temporary files on disk and reads it back in a later phase. ## Where the budget comes from The budget is **per operator instance**, not per query and not per connection. A single query with three sorts and a hash aggregate can claim four budgets at once; run it on 50 connections and the machine can be asked for 200 budgets. Names differ by engine — PostgreSQL `work_mem`, SQL Server's query memory grant, Oracle's PGA workarea target, MySQL's `sort_buffer_size`/`tmp_table_size` — but the shape is the same: a per-operation limit chosen so that concurrency does not blow up total memory use. What counts against the budget is not the table size but the *operator's working set*: for a sort it is the rows being sorted (only the columns carried through the plan, plus per-row overhead); for a hash aggregate it is roughly one entry per **distinct group**, not one per input row. That distinction matters: aggregating 100 million rows into 50 groups needs almost no memory, while aggregating 2 million rows into 2 million distinct keys needs a lot. ## What happens on the disk side For a sort, the engine fills memory, sorts that chunk in place, writes it out as a sorted *run*, and repeats. When input ends it merges the runs. So a spilling sort writes the data once and reads it back at least once — more if there are too many runs to merge in one pass. For a hash aggregate (or hash-based DISTINCT), the engine partitions rows by a hash of the group key. It keeps some partitions in memory and writes the rest to temp files; after the input is consumed it reads each spilled partition back and aggregates it separately. Because all rows with the same key hash to the same partition, this is still correct. Either way, the temporary data lives in a temp tablespace / temp directory, not the buffer pool, and it is not the same thing as normal table I/O. ## Why performance falls off a cliff Three reasons compound: 1. **Extra I/O volume.** Data that was never going to touch storage now gets written and read, often several times its logical size because of per-row overhead in temp format. 2. **A different algorithm.** The in-memory path and the spill path are genuinely different code; the spill path also gives up cache locality. 3. **Shared resource contention.** Temp space is shared. Several concurrent spilling queries fight over the same disk and can push each other over further, producing a feedback loop where a modest traffic increase collapses throughput. Because the switch happens the moment the working set crosses the limit, the curve is a step, not a ramp. This is why "it was fine last month" is such a common story: the table grew 15% and one operator crossed the line. ## How you detect it Every mainstream engine reports it in the executed plan: a sort method of *external merge* with a disk size, a hash aggregate reporting batches/partitions and disk usage, a "temporary table on disk" counter. At the system level you see temp-file creation counts, temp-space bytes written, and I/O on the temp device. Comparing estimated to actual rows on the operator usually explains it too — a wildly under-estimated input means the engine planned for a size that never materialized. ## What to do about it In rough order of leverage: **cut the input** (push filters and projections below the operator so fewer and narrower rows reach it — carrying columns you never use inflates every spilling operator); **avoid the operator** (an index that already delivers the required order lets the engine skip the sort entirely; a top-N pattern avoids sorting the whole set); **fix estimates** so the planner picks the right strategy; and only then **raise the budget**, sized against realistic concurrency rather than for the one query in front of you. Spilling is not a bug — it is the safety valve that keeps a big query from taking the server down. The goal is to make it rare on hot paths, not to eliminate it everywhere.

  • Two queries scan the same number of rows, but only the GROUP BY one spills. Why might that be?
    A hash aggregate's memory need scales with the number of distinct group keys, not with input rows. Grouping 50 million rows into 100 groups keeps a 100-entry table; grouping them by user_id may need millions of entries plus per-group accumulator state. The other query may also be streaming (filter/join with an index) and hold almost no state at all.
  • Is simply raising the per-operation memory limit a safe fix?
    Not globally. The limit is charged per operator instance per query, so a complex plan on many concurrent connections multiplies it; raising it server-wide is a fast route to swapping or an OOM kill. The safer pattern is to raise it narrowly — for a session, role, or batch job that runs at low concurrency — after first trying to reduce the rows and columns reaching the operator.

Sorting a huge pile of paper on a small desk: you sort as much as fits, stack the finished pile on the floor, and repeat — then you still have to walk through all the floor stacks to merge them in order.

saying these in an interview costs you the question

  • Thinking the memory budget is per query or per connection rather than per operator instance
  • Assuming spill risk scales with table size rather than with the operator's working set (rows carried, or distinct groups)
  • Believing temp spill files go through the buffer pool / shared cache like normal table pages
  • Treating spilling as a correctness bug rather than the engine's safety valve
  • Jumping straight to raising the memory setting globally without looking at row estimates or reducing input width

context