A sort in one unit of work wrote 12 GB to the worker's local disk, yet the job finished correctly — why is that by design?
answer
- normal, not a fault
- working set exceeded the borrowed space
- written to disk beside the worker
- same answer, longer runtime
- act only when it repeats
basics
~20 sSpilling is planned degradation: an operator that cannot hold its working set writes part of it to disk attached to the worker and reads it back to finish. The answer is identical; only the runtime grows.
solid answer
~50 sA worker is one process on one machine that owns a fixed amount of memory nothing else can borrow, and several units of work — each one worker's share of one step — run inside it at once and share that memory. When the records a sort must hold at once do not fit in the share it can borrow, it writes part of them out to disk attached to the worker in sorted runs and reads them back before it finishes. That is spilling. The 12 GB is the price paid in disk writes, disk reads and the encoding work on both sides — not evidence of a wrong answer. Most operators in most engines do this rather than fail; some operators in some engines cannot spill at all, and those fail instead. Spilling becomes a problem only when the same records are written and re-read over and over.
go deeper
Recall that an operator with more records than room writes some of them to disk beside the worker and reads them back, and that this is designed behaviour. The answer is the same; the run is slower.
Explain where the bytes go, what the encode and decode on both sides cost, and why seeing bytes on local disk during a step that moves data between workers is not by itself evidence of memory pressure.
Show that you judge spilling by proportion rather than presence: spilled bytes against input bytes per unit of work, and the number of spill events. Reach for dividing the data differently before raising memory everywhere.
The platform decision is the worker shape: total memory against how many units of work run inside one worker. Some spilling everywhere is cheaper than memory reserved and idle on every machine every day.
## The pieces this question turns on A **worker** is one operating-system process on one machine that runs some of the job's work and owns a fixed amount of memory nothing else can borrow. Inside it, several **units of work** — each one worker's share of one step, scheduled and retried on its own — usually run at the same time and share that one budget between them. Part of the budget is **operator working memory**: the space operators borrow while they sort, build a grouping table or build one side of a join, and give back when the step ends. An operator's **working set** is the records it must hold at once to produce a correct result. For a full sort, that is every record in this unit of work's share. When the working set does not fit in the space the operator can borrow, the operator writes part of it out to disk attached to the worker and reads it back before finishing the step. That is **spilling**, and the 12 GB is its trace. ## Why engines are built to do it The alternatives are worse: - **Fail the unit of work whenever the data does not fit.** The working set is not knowable before the run — it depends on how the input divided, how many distinct keys turned up in this share, and how wide the records became once decoded. A system that failed on every misprediction would be unusable. - **Make the author guarantee the fit.** That pushes an unanswerable sizing question onto every job and every input day. Spilling turns a memory shortfall into something that scales gracefully: bytes and seconds. The output is bit-for-bit what it would have been in memory. ## What it costs 1. Encoding the records into the form written out, and decoding them on the way back — often more processor time than the disk transfer itself on a fast local device. 2. The write and the read back. 3. A final pass that reads the written runs together to produce the ordered or grouped result. 4. Space. Disk attached to the worker is finite, and a unit of work that fills it does fail — spilling is graceful until the device is full. ## Three things that are not this | What is happening | What triggers it | Is memory pressure involved | |---|---|---| | An operator spilling its working set | The working set does not fit in the space it can borrow | Yes, that is the definition | | The write side of a redistribution — a step that cannot be computed from records one worker already holds, so every worker writes its output split by destination and every worker fetches its share | Any such wide step | No, it happens whether or not memory was short | | A continuously running job keeping what it remembers per key on the worker's local disk | A deliberate placement choice for that long-lived store | No, it is where that store lives by design | The first is the subject here. The second is the movement subject, and it is the most common confusion: seeing bytes on local disk during a wide step tells you nothing on its own about memory. ## What varies between engines Write no sentence that holds for only one of them: - Some divide a single pool dynamically between operator working memory and **retained-result memory** (the part holding results the job was told to keep for later reuse), so the point at which an operator spills moves with what else is being held. - Some reserve a block the engine manages itself as packed bytes and leave the remainder to the host language runtime, so the spill point is much more predictable but the reserve is fixed. - The oldest model in this class spills out of a sort buffer whose size is set before the job starts and does not move. - Reporting differs too: some engines report spilled bytes per unit of work, others report very little and you infer spilling from step time and disk counters. - Some operators in some engines cannot spill at all. Where that is true, the same shortfall shows up as a failed unit of work instead — a different subject, and a much worse outcome. ## When to stop shrugging Spilling deserves attention when the volume is disproportionate rather than merely present: - spilled bytes several times the unit of work's input bytes; - many spill events inside one unit of work rather than a handful; - step time dominated by disk reads rather than by the computation. And the first responses are not "more memory everywhere": cut the input into more, smaller pieces so each working set is smaller; run fewer units of work concurrently inside one worker so each gets more of the budget; drop columns the step never reads so the records are narrower. Raising memory on every machine is the expensive answer and often the wrong one, because the same total memory divided differently would have fitted. ## Saying it in an interview Name it as designed degradation, name the two costs (disk traffic and the encode/decode on both sides), say the answer is unchanged, and name one signal that would make you act. A candidate who calls every spill a defect has not read a run's own numbers.
- Can spilling ever turn into an outright failure?Yes, in two ways. Disk attached to the worker is finite, so a unit of work can fill it and fail on the write. And some operators in some engines cannot spill at all — where an operator must hold something whole, the shortfall is a failed unit of work rather than slow progress.
- Does spilling change the result in any way?No. The records are read back and processed in full, so the output matches what an in-memory run would produce, including ordering where the step defines one. What changes is elapsed time, disk traffic and processor time spent encoding and decoding the records on the way out and back.
- The step spills a little on most runs. Is that worth tuning?Usually not. A small, stable spill means the budget is close to the working set, which is an efficient place to be — the alternative is memory reserved and unused every day. Act when spilled bytes grow relative to input, when spill events multiply, or when the step starts missing its deadline.
saying these in an interview costs you the question
- Any spilling at all means the job is broken
- Spilled records are dropped, so results are approximate
- Spilling always means the memory budget is undersized
- Spilled bytes are written to the durable shared store the job reads from
- Spilling and the write every wide step performs are the same thing
- More memory on every machine is the only fix