skip to content

A grouping step spills to disk and finishes, but a step that holds one key's whole record set fails — why doesn't disk rescue both?

level: middleimportance: must knowfreq 58%

answer

  1. partial answers can be assembled later
  2. possession versus folding
  3. disk holds parts, never one whole group
  4. what the function retains decides it

basics

~20 s

Writing part of a working set to local disk works only when the answer can be assembled from partial answers. An operation whose contract is possession of one whole thing at one instant has no parts to put aside, so disk has nothing to hold back.

solid answer

~50 s

**Spilling** means an operator that cannot hold its working set writes part of it out — as sorted runs, or as per-key chunks — to disk attached to the worker machine, then reads it back to finish the step. That works when the work can be done on a piece, set aside, and combined later: a sort merges its runs; a fold into a fixed-size accumulator per key can flush part of its **grouping table**, the in-memory table with one entry per distinct key, and merge the partial results afterwards. It does not work when the operation requires one whole thing to exist in memory at a single instant — a group a function takes possession of and indexes or sorts itself, an exact whole-group answer with no partial form, or one individual record that is simply too big. There is nothing to hold back, so more disk and more merging change nothing.

go deeper

for a junior

Recall that writing part of a working set to local disk is designed behaviour, that it makes a step slower rather than broken, and that some operations still cannot be helped by it.

for a middle

Explain the mechanics: runs written and merged, partial per-key values combined, and the contrast with an operation whose contract is to possess one whole thing at a single instant.

for a senior

Demonstrate judgment by naming what varies between engines — streamed group against materialised group — and by reformulating the failing operation into an incremental one instead of buying memory.

for a principal

The angle is which formulations your platform allows by default: whether whole-group possession is a pattern you let teams write at all, and what you offer instead.

## What spilling is, and the assumption underneath it **Spilling** is an operator that cannot hold its working set writing part of that set out — in sorted runs, or in per-key chunks — to **local scratch disk**, meaning disk attached to the worker machine whose contents nobody may read once the job ends, and reading it back to finish the step. It is planned degradation rather than a defect: the step gets slower, sometimes by a large factor, and it completes. It rests on one assumption. The work must be doable on a part of the data, set aside, and combined with the rest afterwards. Every operator that spills successfully has that property; every operator that cannot spill lacks it. That is the whole of the distinction, and it is worth saying in exactly those terms rather than as a list of operator names, because the names differ between engines and the property does not. ## The operations that survive it - **A sort.** Fill memory, sort, write an ordered run, repeat, then merge the runs with a small read buffer per run. Peak memory is governed by how many runs are merged at once, not by the data's size. - **A fold into a fixed-size accumulator per key.** Counts, sums, extremes, and sketch structures all keep one bounded value per key, so the **grouping table** — one entry per distinct key seen so far, growing with distinct keys rather than with records — can be partly flushed and the partial values merged later. - **A join executed as matching chunks.** Split both sides by key into chunk pairs, hold one pair at a time, and the memory needed is one pair rather than one whole side. ## The operations that cannot 1. **A group taken possession of by a function** that builds a list from it, sorts it in place, or indexes it to walk twice. 2. **An exact whole-group answer with no partial form** — an exact median, or an exact enumeration that has to be emitted as a single value. 3. **One individual record larger than the room.** This one is indivisible at any count and under any grouping rule. What these share is that at one instant, one value must exist entire inside one process. The dominant case in practice is the whole group, which is why this leaf is named after it — but the single oversized record fails identically and is not a group, so state the class as "one whole thing at one instant" rather than as groups alone. Cases that look like exceptions to the rule are almost always operators that quietly *do* have a partial form, which is why they survive. | Operation | Has a partial form? | Does writing to disk save it? | |---|---|---| | Sort of a large input | Yes — ordered runs merge | Yes, at a cost in time | | Sum or count per key | Yes — partial values merge | Yes | | Function handed a group it retains | No — it holds the whole set | No | | Exact median of one group | No — needs all values at once | No | | One record bigger than the budget | No — indivisible | No | ## Where engines genuinely differ, and how to state it portably This is the point where a confident answer usually describes one engine and calls it the model. Engines in this class disagree here: - Some **sort the records by key before a user function sees them** and hand that function a lazily advanced iterator backed by the sorted runs. The group itself never has to be resident, and whether the step survives depends entirely on what the *function* retains as it walks. - Others **materialise the group as an in-memory collection** before calling the function at all, in which case the group's own size is the constraint whatever the function then does. - A record-at-a-time runtime does not have this shape at all. What accumulates per key there is long-lived keyed state, kept between records and sized, placed and migrated as its own subject, not an operator's working set that exists only while a step runs. So the portable statement is about retention, not about the engine: whatever the engine's handover style, the engine can only spill structures it allocated and understands. A collection built inside code the engine does not interpret is invisible to its memory manager and cannot be written out on its behalf. ## What to change when you are in the unspillable class - **Replace possession with a fold.** Most whole-group operations have an incremental sibling — a running extreme, a bounded top-N, an approximate quantile with a stated error — and that sibling spills. - **Shrink the resident thing.** Drop fields the operation does not read before the grouping happens; the group that must be held whole is then made of smaller records. - **Group by a finer key that carries the same meaning**, and combine the parts afterwards, so no single resident group is large. - **Raise memory last.** It buys one multiple of headroom, it is paid on every worker in the cluster, and the data's growth takes it back.

  • An aggregation can spill, so why do some aggregations still fail on one unit?
    Because not every aggregation keeps a bounded value per key. A sum keeps one number and merges; an exact median or a collect-the-values operation keeps a structure that grows with the key's record count, and that structure has no partial form to flush.
  • Why does giving the worker more memory not reliably fix the whole-group case?
    It buys one multiple. The resident group's size is set by the data, not by the setting, so a budget doubled today is exceeded by a group that doubles later. It also costs on every worker process in the cluster, for every job that shares the shape.
  • Does a single oversized record behave like an oversized group?
    Worse. A group can sometimes be made smaller by regrouping or by dropping fields; one record is indivisible under every count and every grouping rule. The only remedies are upstream — split the record's payload, or keep the bulk outside the record and carry a reference.

A sink washes a hundred plates a few at a time, which is exactly what spilling is. Seating a hundred people does not work that way: the table has to be big enough at one instant, and taking turns is no substitute for the space.

saying these in an interview costs you the question

  • Believes spilling makes any step fit eventually
  • Says spilling means the job is misconfigured
  • Thinks more local scratch disk removes the limit
  • Assumes every engine materialises a group before the function
  • Treats raising memory everywhere as the first move, not the last