skip to content

Why does an MPP query spill to disk after a shuffle even when total cluster memory looks sufficient?

level: seniorimportance: should knowfreq 48%

answer

  1. the headline RAM number is not one pool
  2. one node's spare memory can't help another node
  3. budget is divided by workers and by concurrency
  4. one worker can receive far more than average
  5. fewer columns crossing means less to hold

basics

~20 s

Cluster memory is not one pool. Each worker gets a fraction of one node's RAM, divided again among concurrent queries, and a shuffle can hand one worker far more than its even share. That worker spills alone while the cluster looks half empty.

solid answer

~50 s

Memory in a shared-nothing cluster is per-node and per-operator, not fungible. A worker's budget is roughly `node RAM ÷ workers per node ÷ concurrent queries`, and the operator above a shuffle — a hash build or a final aggregate — must hold its received partition within that budget. Two things break the arithmetic: **skew**, where one worker receives many times the average partition, and **row width**, where a `SELECT *` carries columns nobody needs above the join. When the budget is exceeded the operator spills partitions to local disk and re-reads them, adding I/O and often extra passes. Spilling to remote object storage instead of local SSD is far worse again. The fixes reduce what crosses the exchange rather than raising limits: project fewer columns, filter earlier, pre-aggregate below the shuffle, fix the skew, or reduce concurrency so each query gets a larger slice.

code

sql · 11 lines
sql
-- Carries every column through the exchange and into the hash table
SELECT *
FROM events e JOIN sessions s ON e.session_id = s.session_id
WHERE e.event_date >= DATE '2026-08-01';

-- Same result set need, far less to move and buffer
SELECT e.session_id, e.event_type, s.country
FROM (SELECT session_id, event_type FROM events
      WHERE event_date >= DATE '2026-08-01') e
JOIN (SELECT session_id, country FROM sessions) s
  ON e.session_id = s.session_id;

go deeper

for a junior

Recall that each machine in the cluster has its own memory and cannot borrow another machine's, so a query can run out of room on one machine while others have space.

for a middle

Explain how a worker's budget is derived — node RAM divided by workers and by concurrent queries — and why the accumulating operator above a shuffle is where the limit is hit.

for a senior

Read a profile and classify the spill: uniform means the intermediate is too big, single-worker means skew, sudden means a plan change. Then reduce shuffled bytes before renting more memory.

for a principal

Own the concurrency-versus-memory trade: admission settings decide how much budget each query gets, so isolation policy and per-query memory behaviour are the same decision, and cluster resizing should be the last lever, not the first.

## Cluster memory is an illusion "The cluster has 2 TB of RAM" is almost never the number that matters. That total is split three ways before an operator sees any of it: 1. **Per node.** Only the memory on the node where a worker runs is usable by that worker; a free gigabyte on another node cannot be borrowed. 2. **Per worker.** Each node runs several parallel workers, so one node's RAM is divided again. 3. **Per concurrent query.** Warehouses admit multiple queries at once, and each is given a share (or competes for a shared pool with a per-query cap). So an operator's real budget can be a small fraction of a percent of the headline number. The relevant question is never "does the dataset fit in the cluster" but "does *this worker's partition* of the intermediate result fit in *this operator's* budget right now." ## Why the shuffle is where it bites Below an exchange, work is local and streaming: a scan reads a block, filters it, and passes it on, holding almost nothing. Above an exchange sits an operator that accumulates — a hash join building its table, a final aggregate holding one entry per group, a sort holding a run. Those are the memory-hungry operators, and they are fed by whatever the exchange delivered. That delivery is where the even-split assumption fails. If the shuffle key is skewed, one worker receives many times the average partition and blows through a budget that was perfectly adequate for its peers. This is the signature MPP memory failure: the profile shows spill on one worker and none on the other 63. ## Row width is half the problem Shuffle and hash-build costs are measured in bytes, not rows. A plan that carries thirty columns through an exchange when the operators above need four is moving and buffering many times the necessary volume. Wide `VARCHAR` columns and semi-structured blobs dominate quickly. Projecting only what is needed — and pushing predicates below the exchange so fewer rows travel at all — often removes the spill entirely without touching a single setting. ## What spilling actually costs When an accumulating operator exceeds its budget, it writes some of its partitions out and continues, then reads them back for a later pass. The costs compound: - **Extra I/O**, both write and read, for data that was already in memory once. - **Extra passes.** If the spilled partitions are themselves too large when re-read, the operator recurses, multiplying the I/O. - **Contention.** Local scratch space is shared with other queries on that node, so one spilling query slows its neighbours. - **Remote spill.** Some cloud engines, having little local disk, spill to object storage. That is orders of magnitude slower than local SSD and turns a degradation into a stall; it is worth knowing whether your engine reports local and remote spill separately, because remote spill is a much louder alarm. Spilling is a graceful degradation, not an error — it is what lets a query bigger than memory complete at all. The failure to avoid is treating any spill as fatal, or treating a large one as normal. ## Diagnosing it Read the profile for the stage that regressed: bytes spilled, per worker. Then ask which of three shapes it is. - **Uniform spill across all workers** means the intermediate is genuinely bigger than the cluster's per-query memory. Options: more memory per query (fewer concurrent queries, or a larger compute cluster), or a rewrite that produces less intermediate data. - **Spill on one worker only** means skew. Fix the key distribution; raising memory only buys time. - **Spill that appeared without a data-volume change** usually means a plan change — an estimate flipped a broadcast into place, or a filter stopped being pushed down after a rewrite. Compare plans, not settings. ## Reducing what crosses the exchange In rough order of leverage: 1. **Project.** Name only the columns needed above the join or aggregate. 2. **Filter early.** Every predicate that can be evaluated before the exchange removes rows from the network and from the hash table. 3. **Pre-aggregate.** A partial aggregation below the shuffle can shrink billions of rows to thousands of partial groups. 4. **Fix skew.** Salt, split, or eliminate hot key values so partitions are even. 5. **Reduce concurrency or resize.** Only after the above — a bigger cluster hides an inefficient plan at recurring cost, and lower concurrency trades one workload's latency for another's. The order matters because the first three shrink the work permanently and for every future run, while the last one rents headroom by the hour.

  • Should you always treat a spill as a failure to be eliminated?
    No. Spilling is the mechanism that lets a query larger than memory finish instead of being killed, and a small spill on a big batch job is often acceptable. What warrants action is a large spill, a spill concentrated on one worker, a spill that recurses into multiple passes, or any spill to remote storage rather than local disk.
  • Why is spilling to object storage so much worse than spilling to local SSD?
    Object storage has far higher latency per request and much lower effective throughput for the small random re-reads a spill produces, and it is charged per request as well as per byte. A join that degrades gracefully when spilling to a local NVMe device can effectively stall when the same partitions have to round-trip to a remote store.
  • Does increasing concurrency ever cause a previously stable query to start spilling?
    Yes, and it is a common surprise. If the engine divides a per-node memory pool among admitted queries, admitting more of them shrinks each one's budget. A report that fit comfortably at concurrency four can spill at concurrency sixteen with identical data, which is why memory-related regressions must be correlated with what else was running.

saying these in an interview costs you the question

  • Adding cluster memory as the first response to any spill
  • Assuming free memory on one node helps a busy node
  • Ignoring column width when reasoning about shuffle size
  • Treating any spill as automatically a failure
  • Not checking whether only one worker spilled

context