What does it mean when an execution plan reports that a sort or hash operation spilled to disk or used temporary storage, and how do you respond?
answer
- blocking operators need a memory grant
- grant sized from estimate x width
- external merge runs / hash partitions to temp
- spill is usually a symptom of misestimate
- memory is per operation x workers x sessions
basics
~20 sThe operator needed more working memory than it was granted, so it wrote partitions or runs to temporary storage and made extra passes. It is usually caused by a low row estimate, an over-wide row, or a memory budget too small for the real data. Fix the estimate first, then the sort or projection.
solid answer
~60 sSorts, hash joins, hash aggregates and deduplication all need a working memory area. The engine sizes it from the estimated row count and width. If the real working set exceeds the grant, the operator degrades to an **external algorithm**: an external merge sort writes sorted runs to temporary files and merges them; a hash operator partitions input to disk and reprocesses partitions one at a time. That is extra I/O plus extra CPU passes, and it is often the difference between a 200 ms and a 40 s query. The plan signals it as spilled/temp bytes, a number of merge passes or hash batches, or an explicit warning. My response order is: (1) check whether the node's actual rows blow past the estimate - a spill is usually a symptom of a cardinality misestimate; (2) remove the operation entirely if an index already supplies the ordering or grouping; (3) narrow the projection, since wide rows and unused columns inflate the working set; (4) only then consider raising the memory budget, keeping per-operation memory times concurrency in mind.
go deeper
Know that sorts and hash operations need memory and fall back to temporary disk files when they do not get enough, which is much slower.
Explain external merge sort and hash partitioning, and that the memory grant is derived from the row estimate.
Give the response order - estimate, eliminate the operator via an ordered access path, narrow the projection, then memory - and mention temp-space contention under concurrency.
Reason about memory as a shared, per-operation budget across concurrency and parallelism, and about when spilling is the correct design for batch work versus a signal of a broken estimate in interactive work.
## Why operators need memory Some operators cannot stream. To sort, you must see all rows before you know which comes first. To build a hash table for a join, you must load the build side before probing. To aggregate without pre-sorted input, you must keep one entry per group. These operators are allocated a **working memory budget** - engines name it differently, but the concept is the same: a per-operation allowance, not a per-query or per-server one. The allowance is sized from the plan's estimated row count multiplied by estimated row width, capped by configuration. That is the first thing to notice: **the memory grant is a function of the estimate**, so a bad estimate causes a bad grant. ## What spilling actually does When the working set exceeds the allowance: - **External merge sort**: the operator sorts what fits, writes that sorted run to temporary storage, and repeats. At the end it merges the runs. If there are more runs than can be merged in one pass, it merges in several passes, and each pass rewrites the whole data volume. - **Hash partitioning (grace hash)**: the operator partitions rows by hash into buckets small enough to fit, writes most buckets to temporary storage, then processes them one at a time. If a partition is still too large - typically because of skew, where one key dominates - it partitions again, and pathological skew can make this very expensive. Costs stack up: temp file writes and reads, extra CPU for repeated passes, cache pollution, and contention on temporary space shared by every session. Under concurrency, several spilling queries can saturate the temp storage and slow down queries that are not spilling at all. ## Reading the signal Plan output for this varies but always includes some of: bytes or blocks written to temporary storage, a count of merge passes or hash batches, peak memory used versus granted, or a categorical marker like 'external' versus 'in-memory'. Two useful readings: - **Spill volume versus data volume.** Spilling 50 MB while processing 40 GB is unremarkable; spilling 40 GB is the query. - **Number of passes/batches.** One pass is a mild degradation. Many passes, or many hash batches with re-partitioning, indicate the grant is drastically undersized or the data is badly skewed. ## Response, in order 1. **Check the estimate at that node.** If actual rows dwarf the estimate, the spill is a symptom; fixing statistics, correlation or an opaque predicate raises the grant automatically and may also change the algorithm to one that does not need the memory at all. 2. **Remove the operation.** Many sorts exist only because the engine had no ordered access path. An index that already provides the ordering (for an ORDER BY, a merge join, or a grouping) lets the operator be replaced by an ordered scan that streams. Likewise a redundant DISTINCT or an unnecessary ORDER BY inside a subquery can often just be deleted. 3. **Shrink the rows.** Working-set size is rows times width. Selecting only the needed columns, avoiding carrying large text or binary columns through a sort, and pushing filters below the blocking operator all reduce the working set - sometimes by an order of magnitude - without touching configuration. 4. **Reduce rows earlier.** Aggregate or filter before joining where semantics allow; a partial/pre-aggregation step can collapse the input dramatically. 5. **Raise the memory budget - carefully.** This is last, not first, because the budget is per operation: a plan with three hash operators and eight parallel workers can consume many multiples of the configured value, times the number of concurrent sessions. Where the engine supports it, raise it for a specific session or a specific reporting workload rather than globally. ## Judgement calls Spilling is not automatically a defect. A nightly batch job that sorts 200 GB will spill by design, and forcing it into memory would be wrong: the right question there is whether the temp storage is fast and whether the pass count is minimal. What matters is whether the spill is *unexpected* - a small interactive query spilling means an estimate or a schema problem - and whether it is *concurrent*, because temp storage and memory are shared resources whose pressure shows up as unrelated queries getting slower.
- Why is raising the working-memory setting globally a risky first fix?Because it is a per-operation allowance, not a per-query or per-server one. A single query with several hash and sort operators, multiplied by parallel workers and by concurrent sessions, can multiply that setting many times over and drive the server into swapping or out-of-memory conditions. Prefer fixing the estimate, removing the sort, or scoping the increase to one session or workload.
- A hash join spills even though its build side is small in total. What would you suspect?Data skew: hash partitioning splits rows by key, so if one key dominates, its partition stays too large no matter how many times it is split, forcing repeated re-partitioning. Look at the distribution of the join key, consider filtering or handling the dominant key separately, or switch the plan toward a merge or index-driven strategy for that key.
saying these in an interview costs you the question
- Treating a spill as a memory-configuration problem before checking the row estimate
- Raising the per-operation memory setting server-wide as the first move
- Assuming any spill means the query is broken, including for large batch jobs
- Not realising SELECT * inflates the working set of every blocking operator
- Believing an index cannot help a sort, only a filter