skip to content

In a columnar warehouse, what makes a window function spill to disk?

level: seniorimportance: should knowfreq 48%

answer

  1. the sort carries every projected column
  2. memory scales with row count times row width
  3. the frame decides how much stays resident
  4. look-ahead frames cannot emit the first row early
  5. MIN over a sliding frame cannot subtract

basics

~20 s

The sort and any row buffering the frame requires. After the shuffle each worker must sort its partitions, carrying every projected column; when a partition plus its frame buffer exceeds the operator's memory budget, the engine writes runs to local disk and merges them back.

solid answer

~50 s

Two things consume memory in a window operator. The **sort** materialises each partition ordered by the window's order key, and it carries the full payload of every projected column, so width multiplies row count. The **frame** then decides how many rows must stay resident: a running aggregate bounded by the current row needs one accumulator and streams; a frame that looks ahead, or a whole-partition aggregate with no ordering, forces the engine to buffer rows before it can emit the first result; a sliding frame keeps the frame's rows, and non-invertible functions such as MIN or MAX cannot cheaply drop the row that leaves. When the total exceeds the budget, the engine writes sorted runs to local disk and merges — correct, but with extra I/O and a much longer stage. The levers are: fewer rows into the window, fewer columns carried through it, smaller partitions, a cheaper frame shape, or more memory per worker.

code

sql · 9 lines
sql
-- Before: the sort carries every column of a wide fact table
SELECT *,
       sum(amount) OVER (PARTITION BY account_id ORDER BY event_ts) AS running
FROM ledger_wide;

-- After: only the columns the report needs travel through exchange and sort
SELECT account_id, event_ts, amount, currency,
       sum(amount) OVER (PARTITION BY account_id ORDER BY event_ts) AS running
FROM ledger_wide;

go deeper

for a junior

Know that a window function has to sort rows, that sorting needs memory, and that when memory runs out the engine writes to disk and slows down. Selecting only the columns you need is the easy first improvement.

for a middle

Separate the two consumers: the sort, whose cost scales with rows times row width, and the frame, which decides how many rows stay resident. Explain why a running total streams while a look-ahead or whole-partition frame must buffer.

for a senior

Diagnose from a profile — per-worker peak memory and spill bytes on the window stage — and apply the fix ladder in order: fewer rows, fewer columns, smaller partitions, cheaper frame shape, then memory. Note that spill is per worker and skew concentrates it.

for a principal

Decide where spill is worth paying for. A nightly batch that spills inside its window needs no intervention; an interactive dashboard shape that spills on every refresh justifies changing the physical layout or precomputing, and that call belongs at the platform level.

## Where the memory goes A window operator in a columnar MPP engine runs in three phases: receive rows from the repartitioning exchange, sort them by (partition key, order key), then scan the sorted stream emitting one result per row. Memory is consumed in the second and third phases, and it is worth separating them because the fixes differ. **The sort.** Sorting is not free-standing — it carries the payload. Every column the query still needs downstream travels with the sort key, so a `SELECT *` over a 200-column fact table sorts vastly more bytes than a query that projects the six columns it actually uses. Columnar engines sort batches of column vectors, but the total resident bytes still scale with row count times row width. When that exceeds the operator's budget, an external merge sort kicks in: sorted runs are written to local disk and merged back on the read side. **The frame buffer.** After sorting, how much has to stay in memory depends entirely on the shape of the frame: - A frame that ends at the current row and starts at the beginning of the partition — the classic running total — needs a single accumulator. Constant memory, fully streaming, first row out immediately. - A frame that extends *ahead* of the current row, or that spans the whole partition with no ordering, cannot produce the first row until later rows have been read. The engine must buffer forward, in the worst case the entire partition. - A sliding frame of fixed width keeps the rows inside the frame. For invertible aggregates — SUM, COUNT, AVG — the engine can add the entering row and subtract the leaving one, so the state stays small even for wide frames. For non-invertible ones — MIN, MAX, and anything distinct-based — removing a row can change the answer arbitrarily, so the implementation must retain the frame's values (or use a specialised structure), and cost grows with frame width. - Frames defined by value ranges rather than row offsets must materialise **peer groups**: all rows sharing the same order-key value. If the order key is low-cardinality — an hour truncated timestamp, a date, a status — a single peer group can be enormous, and this surprises people who assumed a range frame behaves like a row frame. ## Skew makes it worse Spill is per worker, and workers do not hold equal shares when the partition key is skewed. A cluster with plenty of aggregate memory can still spill on one node because that node owns the one partition that is a hundred times bigger than the median. Always look at the *maximum* per-worker memory and spill in the profile, not the average. ## Recognising it Runtime profiles report bytes written to local storage, or a spill/temporary-storage indicator, on the sort or window operator. The symptoms are a stage whose duration jumps disproportionately with a modest increase in data volume, high local disk I/O, and (in cloud engines that spill to remote storage when local disk fills) a second, far steeper cliff. Spilling is a correctness-preserving fallback, not an error, so nothing fails — the query just gets slow and expensive. ## The fix ladder 1. **Fewer rows.** Push filters below the window. Rows removed by a predicate applied after the window were still shuffled, sorted and spilled. Where the report allows it, pre-aggregate to a coarser grain first: a window over daily totals per account is orders of magnitude smaller than one over raw events. 2. **Fewer columns.** Replace `SELECT *` with the columns the result needs. This is often the largest single win in a wide fact table and costs nothing semantically. 3. **Smaller partitions.** Adding a secondary partition column, where the semantics allow, both spreads work across workers and reduces the biggest partition — the one that determines whether anything spills at all. 4. **A cheaper frame.** Ask whether the query really needs a look-ahead frame or a range frame. Rewriting a value-range frame as a row-offset frame, when the data has one row per period, removes peer-group materialisation entirely. 5. **Exploit existing order.** If the table's physical layout already sorts rows by the window's order key inside each partition, some engines can stream the window without a full sort. That is a storage-layout decision that pays off for every query using the same shape. 6. **More memory.** Raising the per-operator budget or moving to workers with more memory is the last lever, and the honest answer when the data volume is simply large and the query runs rarely. ## When spill is acceptable A batch job that spills once a night and finishes inside its window is fine — buying memory to avoid a disk write nobody waits on is wasted money. Spill matters when it lands in an interactive path, when it happens on every dashboard refresh, or when it turns a linear cost curve into a cliff as data grows. Decide by where the query sits, not by the presence of the word "spill" in a profile.

  • Why does a sliding SUM stay cheap over a wide frame while a sliding MIN does not?
    SUM is invertible: as the frame slides the engine adds the entering value and subtracts the leaving one, keeping constant-size state. MIN cannot be undone — removing the current minimum leaves no way to recover the next one from an accumulator — so the implementation must retain the frame's values or maintain a specialised structure, and its cost grows with frame width.
  • A query spills only on one worker even though the cluster has ample total memory. Why?
    Memory budgets are per worker and per operator, and the exchange does not distribute rows evenly when the partition key is skewed. The worker owning the largest partition can exceed its own budget while peers use a fraction of theirs. Aggregate memory is irrelevant; the maximum per-worker partition size is what determines whether anything spills.
  • Does replacing SELECT * with an explicit column list really change window memory?
    Yes, and often more than any other single change. The sort materialises rows with the full payload of every projected column, so a 200-column fact table narrowed to six columns can shrink the sorted bytes by an order of magnitude, cutting both the network transfer through the exchange and the chance of spilling.

saying these in an interview costs you the question

  • Treats spill as an error rather than a fallback that trades memory for I/O
  • Thinks every window frame buffers the whole partition
  • Compares total cluster memory instead of per-worker budgets
  • Believes wider sliding frames always cost more for every function
  • Adds memory before removing unused columns and rows

context