skip to content

A unit of work that overflowed to local disk reports 9 GB spilled in memory but 1.1 GB on disk — why both?

level: middleimportance: nice to knowfreq 35%

answer

  1. same records, two forms
  2. as held versus as encoded
  3. encoding and compression explain the gap
  4. size memory from the larger figure
  5. read it per unit of work

basics

~20 s

Both figures describe the same records in two forms: the space they occupied as the operator held them, and the bytes written after the engine encoded them. The gap is a fact about representation, not lost or duplicated data.

solid answer

~40 s

Spilling is an operator writing part of its working set to disk attached to the worker and reading it back. Engines commonly report it twice. The larger figure is the size those records occupied in the space the operator had borrowed — as ordinary objects of the worker's language, each carrying the runtime's own per-object bookkeeping, or as the engine's own packed layout. The smaller figure is the bytes actually written, after encoding into a compact form and sometimes compressing on top. Nothing is lost between them. Use the in-memory figure to answer "how much more operator working memory would have avoided this", and the on-disk figure to answer "how much disk traffic did I actually pay for". Answering the first question with the second is the common error, and it undersizes memory badly.

go deeper

for a junior

Recall that a spill can be reported twice: the space the records took in memory and the bytes written out. They describe one event in two forms, and nothing is lost between them.

for a middle

Explain the gap through representation: per-object bookkeeping or a packed layout in memory against a compact encoded and possibly compressed form on disk. Say which figure answers a memory question.

for a senior

Read the pair per unit of work, not as a step total, and use the spread across units to tell a genuine shortfall from one unit holding far more of the data than the rest.

for a principal

The ratio is a platform signal. If it is consistently large across jobs, how records are held is costing the organisation memory everywhere, and that is a default worth changing rather than a per-job fix.

## Two numbers, one event When an operator cannot hold its working set and writes part of it out to disk attached to the worker, a run's reported numbers frequently show that single event twice, under two names and two very different sizes. Neither is wrong and neither is a duplicate count. - **The in-memory figure** is how much of the operator's borrowed space those records were occupying at the moment they were evicted from it. - **The on-disk figure** is how many bytes were actually written to local scratch disk. The records are the same records. What differs is the form they are in. ## Why the forms differ so much A record in an operator's working memory is held in one of two ways: - **As host-runtime objects** — ordinary objects of whatever language the worker runs, each carrying the runtime's own per-object bookkeeping, references between them, and padding. A small logical record can occupy several times its logical size this way. - **As engine-managed binary records** — a packed byte layout the engine allocates and interprets itself, where a field is read at a known offset instead of by following a pointer. This is much closer to the logical size, though still not equal to it. What is written to disk is neither: it is a compact encoded form, often with a compression step on top, chosen because disk traffic is the thing being minimised. Column-shaped or dictionary-friendly data can shrink dramatically; opaque binary payloads barely shrink at all. So the ratio is a property of **the data and the engine's representation**, not a bug. It is often several times, and much smaller on engines that already hold their records as packed bytes. | Figure | What it measures | Chiefly tells you | |---|---|---| | In-memory size | Space the records held in the operator's borrowed memory | How much more operator working memory would have avoided the spill | | Bytes on disk | Bytes written to local scratch disk | The disk traffic and space cost you actually paid | | Their ratio | The encoding and compression multiple for this data | Whether the records are representation-heavy or already compact | ## Which one answers which question This is where candidates lose the point. The instinct is to take the smaller number, because it looks concrete, and to conclude that a little more memory would have avoided the whole thing. That reasoning is backwards: to avoid spilling, the records would have had to stay **in the form they had in memory**, which is the larger number. Sizing memory from the on-disk figure undersizes it by exactly the encoding multiple, and produces a job that still spills after the change. Conversely, when the question is cost — the time this step spent writing and reading, and whether the local device has space — the on-disk figure is the honest one. Charging the in-memory size against disk would overstate it just as badly. A third reading comes from the pair together. An unusually large ratio says the records are expensive to hold relative to what they really contain, which points at representation rather than at volume; an unusually small ratio says the data is already dense and there is no encoding win left to find. ## What varies between engines - Some engines report both figures per unit of work; some report only one; some aggregate to the step and leave you unable to see whether one unit of work produced nearly all of it. - Engines that hold packed bytes show a smaller gap than engines that hold ordinary language objects, so a ratio that is remarkable on one platform is unremarkable on another. - Where the numbers are aggregated across a step, a large total can mean either every unit of work spilled a little or one spilled enormously — two situations with different remedies, which is why the per-unit view matters more than the total. - Some engines count the final merging pass's reads in the same counter, and some do not, so comparing the ratio across platforms is not meaningful. ## How to use it in practice 1. Read the **per unit of work** figures, not just the step total, and look at the spread across units. 2. Take the in-memory figure as the memory question and the disk figure as the cost question. 3. Before asking for more memory everywhere, check whether more, smaller pieces would bring the in-memory figure under the space each unit of work can borrow — the same total memory, divided differently. 4. Treat a large ratio as a hint about how records are held, which is its own subject, rather than as an anomaly in the spill itself. The short version for an interview: the two numbers are the same records in two forms; size memory from the memory one and cost from the disk one.

  • Which figure would you use to decide how much more memory to give the step?
    The in-memory figure. Avoiding the spill means keeping those records in the form they had in memory, so that is the space needed. Sizing from the bytes on disk undersizes by the whole encoding multiple and leaves the step spilling after the change.
  • What does an unusually small gap between the two figures suggest?
    That the records are already dense — either the engine holds them as packed bytes rather than as language objects, or the payloads resist compression. There is little encoding win left, so the remedies are volume-shaped: fewer records per unit of work, or narrower records.
  • The step total is large but every unit of work looks similar. Does that change your response?
    Yes. An even spread means the working set genuinely exceeds what each unit can borrow, so dividing into more pieces or narrowing records helps. A total dominated by one unit points instead at an uneven share of the data, which is a different subject with different remedies.

saying these in an interview costs you the question

  • The two figures mean the spill was counted twice
  • Size memory from the bytes written to disk
  • The smaller figure proves data was dropped
  • The ratio is the same on every engine and every dataset
  • Only the step total matters, not the per-unit figures