When an operator runs out of room, what does a sort write to local disk, and how does a grouping table differ?
answer
- records against distinct keys
- ordered runs versus flushed partial entries
- final pass takes the smallest next
- partials must merge pairwise
- one buffer per run at the end
basics
~20 sA sort writes ordered runs and later reads them together, always taking the smallest next record. A grouping table instead flushes the partial entries built so far, starts fresh, and combines partials for the same key at the end.
solid answer
~50 sBoth are spilling — an operator writing part of its working set out to disk attached to the worker and reading it back — but the shapes differ. A sort fills the space it can borrow, orders what it holds, writes that as one ordered run, and repeats; at the end it reads the runs together, repeatedly taking the smallest next record across them, which needs only a small buffer per run. A grouping table holds one entry per distinct key seen so far, so it grows with distinct keys rather than with records; when it is full it writes the entries it has as partial results, starts an empty table, and at the end merges partials for the same key. That last step is cheap when partial results can be merged pairwise, such as a running sum and count, and impossible when the aggregation must see every value for a key at once.
go deeper
Recall the two shapes: an ordered run written out and read back together with others, and a table of partial results per key flushed and merged later. Both let a step finish on more data than fits.
Explain why the final pass of a sort needs only a buffer per run, and why a grouping table's size follows distinct keys. Name an aggregation whose partials merge pairwise and one whose do not.
Diagnose from the shape: a sort spilling means too many records per unit of work, a grouping table spilling means too many distinct keys. Those point at different remedies, and memory is rarely the first one.
The platform question is which aggregations authors are encouraged to write. A house habit of fixed-size, pairwise-mergeable accumulators keeps grouping steps spillable and predictable at any key count.
## Two operators, two working sets Spilling means an operator that cannot hold its working set writing part of it out to disk attached to the worker and reading it back to finish the step. What it writes depends entirely on what the operator was holding and why. A **sort** must hold every record in this unit of work's share before it can emit the first one in order. Its working set grows with the **number of records**. A **grouping table** — the in-memory table an aggregation builds, one entry per distinct key seen so far — only has to hold one accumulator per key. Its working set grows with the **number of distinct keys**, not with the number of records. A billion records over fifty keys is nothing; ten million records over ten million keys is the same size as the input. ## How a sort spills 1. Fill the space the operator can borrow with incoming records. 2. Order what is held and write it out as one **ordered run** on local scratch disk. 3. Empty the space and carry on, producing run after run. 4. At the end, read the runs **together**: keep one record from each in memory and repeatedly emit the smallest, refilling from whichever run it came from. Step 4 is the important one. The final pass does not need to hold all the data — only a read buffer per run and one record from each. That is why sorting far more data than memory is an ordinary operation rather than an impossible one, and why a sort that spills usually writes each record once and reads it once. ## How a grouping table spills 1. Build entries, one per distinct key, updating the accumulator for a key already present. 2. When the table no longer fits, write the entries it currently holds as **partial results** for those keys, then start an empty table. 3. Carry on. The same key can now appear in several flushed batches, because after a flush the table no longer knows it was seen. 4. At the end, bring together the partials for each key and apply a merging function to them. The difference that matters is step 4. It is cheap when two partial results for a key can be merged into one — a running sum, a count, a minimum, a sum and a count carried together for a mean. It is expensive when the accumulator is large, and it does not work at all when the aggregation needs every value for a key present at once, such as an exact median over the whole group. Where an aggregation cannot be expressed as merged partials, the operator has nothing useful to flush, and the shortfall becomes a failure rather than a slowdown — that is a separate subject from spilling. | | A sort | A grouping table | |---|---|---| | Working set grows with | records in this unit of work's share | distinct keys seen | | What is written out | ordered runs | partial results for the keys held so far | | What the final pass does | reads runs together, taking the smallest next record | merges partials for the same key | | Cheap when | always, given enough room for one buffer per run | partial results merge pairwise | | Badly behaved when | the room per run buffer is too small for one pass | the aggregation needs all values for a key at once | ## What varies between engines The shapes above are the class's shapes, but the trigger and the accounting differ: - Where a single pool is divided dynamically between operator working memory and results the job was told to keep, the point at which either operator spills moves during the run, depending on what else is held. - Where the engine reserves a block it manages itself as packed bytes, the spill point is stable and the records written out are already in a compact layout. - Where the operator spills out of a buffer whose size is fixed before the job starts, the spill point does not move at all and is the author's to set. - Some engines partition the grouping table by key before flushing so that each key's partials land together; others sort the flushed entries instead. Both end at the same place by different routes. - A continuously running job usually aggregates into what it remembers per key between records rather than into a transient grouping table — a different store with different sizing, owned by the state subject, not this one. ## The practical reading When a sort spills, expect roughly one write and one read of the data, and suspect trouble when the volume read back is several times what was written. When a grouping table spills, ask about distinct keys first: adding memory helps for a while, whereas an aggregation whose partials merge pairwise stays comfortable at almost any key count, because it never has to hold more than one accumulator per key per flush.
- Why does the final pass of a spilled sort not need memory proportional to the data?Because each run on disk is already ordered. The pass only has to see one record from each run at a time and take the smallest, refilling from that run. Memory scales with the number of runs and the read buffer per run, not with total records.
- Which aggregations keep a grouping table small regardless of record count?Those whose accumulator is fixed-size per key and whose partials merge pairwise — counts, sums, minima and maxima, a sum-and-count pair for a mean, and bounded sketch-style summaries. The table then grows only with distinct keys, so record volume is almost irrelevant to its footprint.
- Why can the same key appear in several flushed batches of a grouping table?Because a flush empties the table. Records for a key seen before the flush contributed to a partial that is now on disk; records for the same key arriving afterwards start a new entry. Only the final merging step brings those partials back together.
- Does a grouping table's spill preserve the aggregation's exactness?Yes, where the aggregation's partials merge exactly — sums, counts, minima. A summary that was approximate to begin with stays approximate to the same degree. Spilling introduces no error of its own; it only changes when and how the accumulators are combined.
Sorting a huge stack of paper on a small desk: fill the desk, order that batch, put the ordered pile on the floor, repeat. At the end, keep the top sheet of every pile visible and keep taking the lowest one. You never need a desk big enough for all the paper — only one big enough to see one sheet from every pile at once.
saying these in an interview costs you the question
- A spilled sort must re-read everything into memory at the end
- A grouping table grows with records rather than distinct keys
- Flushed partial groups for one key are simply overwritten
- Any aggregation can be spilled and merged later
- Sorting more data than memory is impossible