One unit of work writes hundreds of small sorted runs to local disk and reads several times its input back — when has spilling stopped being acceptable?
answer
- written once, read once, roughly
- tiny runs mean many rounds
- bytes read against input bytes
- budget divided by concurrent units
- more pieces before more memory
basics
~20 sSpilling stops being acceptable when the same records are written and re-read repeatedly. Healthy spilling moves the data out and back about once; hundreds of tiny runs mean the room per unit of work is too small for even one combining pass.
solid answer
~50 sHealthy spilling has a bounded shape: an operator writes part of its working set out to disk attached to the worker, then reads it back once to finish. The cost is roughly one write and one read of the data. Thrashing is the same mechanism run past its useful range: the space one unit of work can borrow is so small relative to its share that it produces very many tiny runs, and the final pass cannot keep even a small read buffer for each of them at once, so the runs must be combined in several rounds — each round rewriting and re-reading the data. Bytes moved then grow with the number of rounds, not with the data. The signal is the ratio of bytes read back to input bytes per unit of work, together with the run count. The remedies are usually more, smaller pieces, or fewer units of work running concurrently inside one worker.
go deeper
Recall that writing to disk once and reading it back once is the normal shape. Doing it over and over is not, and hundreds of tiny written runs from one unit of work is the visible sign.
Explain why tiny runs force several rounds of combining, and that each round rewrites and re-reads the data. Name the ratio you would check: bytes read back against input bytes per unit of work.
Demonstrate the concurrency divisor: room per operator is the worker's operator working memory divided by the units running at once. Check the per-unit spread before proposing any change, and prefer dividing the data differently.
The standing decision is worker shape: how much memory against how many concurrent units of work, set as a default that every job inherits. Set it too aggressively and whole classes of steps thrash under load.
## Where the line falls 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. In its healthy form the accounting is simple: **each record is written once and read once**, plus whatever the final pass costs to produce the ordered or grouped result. Runtime grows by a predictable amount, and the answer is unchanged. Thrashing is not a different mechanism. It is the same one operating below the point where one pass is possible: 1. The space a unit of work can borrow is small relative to its share of the data, so each ordered run written out is tiny. 2. Tiny runs mean very many runs. 3. The final pass needs a read buffer for each run it reads at once. If there is not room for all of them, the runs must be combined in **rounds** — read some, write a larger merged run, repeat. 4. Each round writes and re-reads the whole data set again. So the data volume moved is no longer about one write and one read; it multiplies by the number of rounds, which grows as the space per unit of work shrinks. A step that would have taken minutes takes an hour, and the step time is dominated by disk reads rather than by the computation. ## The arithmetic that actually bites The number people forget is **concurrency inside one worker**. A worker is one process owning 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. The space one operator can borrow is roughly the worker's operator working memory divided by the units of work running concurrently. That has an uncomfortable consequence: **raising a worker's memory and its concurrency in the same proportion changes nothing per operator**. Twice the memory and twice the threads is the same room each, on twice as much data per worker. Teams do this routinely and report that more memory did not help — and they are right, because per-operator room never moved. ## Reading the run's own numbers Per unit of work, not per step: - **Bytes read back against input bytes.** Around one is healthy. Several times over is the thrashing signature. - **Number of spill events or runs.** A handful is fine. Hundreds from one unit of work says each run is tiny. - **Time split.** A step whose elapsed time is mostly disk reads while processor usage is low is doing transport, not work. - **Spread across units.** If one unit of work spills hugely and its peers barely spill, the problem is that its share of the data is much larger than the rest — an uneven-share problem whose diagnosis and remedies belong to the skew subject, not to this one. ## What to change, in order 1. **Cut the input into more, smaller pieces.** Each unit of work then holds less, so the working set may fit or spill in one pass. This is free in the sense that no machine changes; it is the same total memory divided differently. Taken too far it produces scheduling overhead and tiny outputs, so it is not unbounded. 2. **Run fewer units of work concurrently inside one worker.** This trades parallelism for room per operator, and it is the right trade when the alternative is several rounds of re-reading. Elapsed time often falls even though fewer things run at once. 3. **Narrow the records.** Drop columns the step never reads before the sorting or grouping operator sees them; the working set shrinks in direct proportion. 4. **Then, and usually only then, more memory per worker** — keeping concurrency fixed, or the change cancels itself out. ## What varies between engines - Where a single pool is divided dynamically between operator working memory and results the job was told to keep, a job holding many kept results squeezes operators into thrashing without the author changing anything about the step itself. - Where the engine reserves a block it manages as packed bytes, the room per operator is more predictable and records are denser, so the same data produces fewer, larger runs. - Where the operator spills out of a buffer whose size is fixed before the job starts, the run count is almost directly the author's choice, and a badly chosen value produces thrashing deterministically on every run. - Some engines adjust the number of pieces at runtime from what they measured, which can relieve this by itself; others do not, and in a continuously running job the width is usually fixed for the life of the job. Do not assume the engine will notice. ## The judgment to show A senior answer distinguishes **presence** from **proportion**, names the concurrency divisor rather than only the worker total, reads the per-unit spread before proposing anything, and reaches for dividing the data differently before asking for memory on every machine. A candidate who answers "raise the memory" has skipped every measurement that would tell them whether it will work.
- Why does raising worker memory sometimes change nothing at all?Because the room one operator can borrow is roughly the worker's operator working memory divided by the units of work running inside it at once. Raising memory and concurrency together leaves that quotient unchanged, on more data per worker. Hold concurrency fixed, or the increase cancels itself.
- What single ratio would you check first on a step you suspect is thrashing?Bytes read back from local disk against input bytes, per unit of work. Around one means the data moved out and back once, which is healthy. Several times over means it is being re-read in rounds, and the run count will confirm it.
- One unit of work spills enormously and the rest barely spill — is that thrashing?Not the same problem. An uneven spread points at that unit holding a much larger share of the data, so the fix is about how the data divides rather than about room per operator. Its diagnosis and remedies belong to the uneven-work subject.
- Can reducing concurrency inside a worker make a step finish faster?Yes, when the step is thrashing. Fewer units of work at once means more room each, which can collapse several rounds of re-reading into one pass. Less runs in parallel but each finishes far sooner, so elapsed time falls. It helps only while thrashing is the constraint.
saying these in an interview costs you the question
- Any amount of spilling is equally acceptable
- More worker memory always relieves spilling
- Raising memory and thread count together doubles room per operator
- The step total is enough; per-unit figures add nothing
- The engine will always notice and replan around it
- Thrashing and an uneven share of the data are the same problem