skip to content

Memory, Spill and Caching

What happens when a worker runs out of room: how one process's memory is divided, what it writes to local disk instead, and which failures more memory will never fix.

on this pageshow

explore

questions

20

Two later steps each read the same computed intermediate — what does pinning that result change about the work done?

level: juniorimportance: must knowfreq 60%

answer

  1. count the readers first
  2. nothing is kept by default
  3. each reader walks the branch again
  4. one reader buys you nothing
  5. kept copy sits in retained-result memory

basics

~20 s

Pinning tells the engine to keep the computed intermediate after the first reader finishes, so the second reader reads the kept copy instead of re-running every step that produced it. Without a pin, that branch is computed twice.

solid answer

~50 s

In most engines of this class a named intermediate is a description, not a stored thing: records flow through the operators and are dropped as they go. So each later step that asks for the intermediate walks the whole branch above it again, back to the source, and repeats every step in it. Pinning a result tells the engine to keep the computed intermediate — each worker keeps its own share — so the second and later readers read the kept copy instead. The saving is `n - 1` walks of the branch where `n` later steps read it, and exactly nothing where only one does. What "kept" means varies: some engines hold it as ordinary objects of the worker's language, some as a packed byte layout, some overflow it to disk attached to the worker, and on several a pin is a request the engine may drop under memory pressure rather than a guarantee.

go deeper

for a junior

Recall that a named intermediate is normally recomputed every time a later step reads it, and that pinning keeps the computed copy so the second reader skips that work. Say out loud that one reader means no saving.

for a middle

Explain where the kept copy lives — the worker's retained-result share, next to the operator working memory the next step needs — and that engines differ over whether a pin is a guarantee or a request they may drop.

for a senior

Show that you decide by counting readers and by reading the run's own numbers afterwards, and that you check whether the steps after the pin started spilling, because that is where a pin turns into a loss.

for a principal

The call you own is whether authors may pin at all on a shared platform, what the default form is, and who is accountable for releasing pins — because an unreleased pin costs every other job on the same worker.

## What "a computed result" means here A **computed result** is a named intermediate your program describes — the output of a chain of steps over some input — that later steps read. In most engines of this class the program is a *description* first: you name the steps, and work begins when something demands an answer (a count, a write, a sample pulled back). When that demand arrives, the engine walks back up the **branch above** it — each step pulling from the step before, up to the source — and runs the whole chain. Nothing in that walk is kept by default. Records flow through the operators, are consumed, and are dropped. That is deliberate, not an oversight: the alternative is holding every intermediate of every step, and no worker has room for that. The consequence is the one this question turns on — a **second** demand that reads the same named intermediate walks the same branch again and re-runs every step in it. ## What pinning does **Pinning a result** means telling the engine to keep a computed intermediate so later steps read it instead of recomputing the branch that produced it. The work is distributed, so the copy is too: each **worker** — one operating-system process on one machine that runs part of the job and owns a fixed amount of memory nothing else can borrow — keeps its own share of the result, next to the units of work running inside it. The copy is held in **retained-result memory**: the part of a worker's budget holding results the job was told to keep, as distinct from **operator working memory**, the part operators borrow while sorting, building a grouping table or building the held side of a join and give back when the step ends. What "kept" physically means is not one thing: | where the copy can live | cost to read it | footprint | |---|---|---| | as host-runtime objects — records as ordinary objects of the worker's language | cheapest read | largest; each object carries the runtime's own bookkeeping | | as a packed byte layout the engine lays out and interprets itself | a decode step per read | markedly smaller | | on local scratch disk attached to the worker | a disk read per use | bounded by disk rather than memory | | nowhere, because the engine dropped it under memory pressure | a full recompute of the branch | none | That last row is the one beginners miss. On several engines a pin is a **request**, not a promise: when operator work needs room the engine may perform *the dropping of a pinned result to make room for operator work*, and the next reader quietly recomputes the branch. ## The one thing that justifies a pin Reuse, and only reuse. Count the later steps that actually read the intermediate: 1. **One reader** — the branch runs exactly once either way. The pin buys nothing and still holds memory. 2. **Two or more readers** — the branch runs once instead of `n` times. The pin buys `n - 1` walks of it. 3. **Pinned and never read** — pure loss: memory held for a result nobody asks for. The honest cases are concrete: two outputs written from the same filtered and cleaned set; one prepared set feeding both a join and an aggregate; an iterative computation that reads the same prepared set every round. ## What varies between engines - Engines built around finite input generally expose an explicit pin. A **record-at-a-time continuous runtime** largely does not: it has no finite computed result to hold, and what it remembers between records is long-lived keyed state, which is a different mechanism with a store of its own. - The oldest two-phase disk-to-disk model in this family writes its intermediate out between phases whether or not you asked, so "reuse" there means reading a written output again rather than pinning anything. - An engine that runs continuous work as a rapid succession of small finite jobs sits in between: a pin can hold within one of those small jobs and not across them. ## How you check it worked Settle it from **the run's reported numbers** — whatever the engine reports per unit of work after a run — rather than from arithmetic: - the second reader's steps should show the kept copy being read, not the branch's steps appearing a second time; - the held bytes should be visible somewhere, in a form you can reason about; - spill figures on the steps that follow should not have grown. If they have, the pinned copy is displacing operator working memory, and the pin is a trade rather than a free win. ## One thing a pin is not A pin is not a durability device. A consistent picture written so a crashed job can resume is a **recovery snapshot** — a different mechanism, written to durable shared storage, whose purpose is surviving failure rather than saving recomputation. Treating one as the other is the most common confusion on this subject.

  • What happens when a pinned result does not fit in the worker's retained-result memory?
    It depends on the engine and on the form you asked for. Some keep what fits and recompute the rest of the branch on demand. Some overflow the surplus to local scratch disk attached to the worker, trading a disk read for a recompute. Some hold it in a packed byte layout that is several times smaller than ordinary objects, at a decode cost per read. The failure mode to watch for is paying twice: holding bytes and still recomputing.
  • Why does pinning an intermediate that only one later step reads still cost something?
    Because the memory the copy occupies comes out of the same worker process that operators are drawing on. Retained-result memory and operator working memory are shares of one fixed budget, so a copy nobody reads twice still shrinks the room available for the next sort, grouping table or held join side, and can push that step into writing part of its working set out to local disk.
  • Does pinning change the answer the job produces?
    It should not, and where it does you have found a real problem: a branch whose result differs between two walks. That happens when the branch reads a source that changes under it, or contains something time-dependent or random. In that situation the pin is not an optimisation — it is quietly choosing one of the two answers and freezing it, which is worth noticing before you rely on it.

saying these in an interview costs you the question

  • Pinning always makes a job faster, so pin every intermediate
  • A pinned copy is free when the cluster has spare capacity
  • Pinning a branch that only one step reads still saves that step time
  • Once pinned, a result is guaranteed to stay in memory
  • Pinning is how a job survives a worker dying
open as a page

One unit of work fails for lack of memory four times and kills the job while the cluster sits idle — what does that rule out?

level: juniorimportance: must knowfreq 70%

basics

~20 s

Memory is owned per worker process, so a failure inside one unit of work says nothing about total cluster capacity. Idle machines cannot lend memory to the process that is short; what reached that one unit has to change instead.

open as a page

Why can one record held as an ordinary language object occupy several times the memory of the same record packed as bytes?

level: juniorimportance: must knowfreq 62%

basics

~20 s

Each field can become a separate object carrying the runtime's own per-object bookkeeping - a header, alignment padding, and a pointer to it from the record. A record of many small fields therefore costs a multiple of its raw bytes, not a small addition.

open as a page

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?

level: juniorimportance: must knowfreq 62%

basics

~20 s

Spilling 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.

open as a page

A worker process is given one fixed memory budget. What competing demands share it while the job runs?

level: juniorimportance: must knowfreq 70%

basics

~20 s

Operators borrowing room while they run, any results the job asked the engine to keep, and runtime overhead the engine never counts - all shared by every unit of work inside that one process, with no lending from idle workers.

open as a page

After a large intermediate is pinned, a job slows on unchanged input — which memory did the copy take, and from what?

level: middleimportance: must knowfreq 55%

basics

~20 s

The pinned copy sits in the same worker process whose operators need room to sort, group and build join sides. A large pin shrinks that working memory, so later steps write part of their working set to disk and the job slows.

open as a page

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%

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.

open as a page

A worker is killed outright while the engine's own numbers show its budget well below full. What explains that?

level: seniorimportance: must knowfreq 55%

basics

~20 s

The engine reports only the memory it accounts for. The platform's ceiling covers the whole process, including the language runtime, buffers and user-function allocations the engine never counted - so the process crosses a line the engine could not see.

open as a page

Why does a worker process allocating an object per record show repeated pauses rather than a memory error?

level: middleimportance: should knowfreq 48%

basics

~20 s

Because the memory is being freed successfully. Where the worker's language runtime reclaims automatically, it periodically finds and frees objects nothing refers to, stopping the worker's own work while it does. Huge allocation volume buys time spent reclaiming, not a refused allocation.

open as a page

When an operator runs out of room, what does a sort write to local disk, and how does a grouping table differ?

level: middleimportance: should knowfreq 55%

basics

~20 s

A 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.

open as a page

An engineer sizes a worker's memory by multiplying input bytes by a guessed factor. What should replace that arithmetic?

level: middleimportance: should knowfreq 50%

basics

~20 s

A measured run. Read what the engine reported per unit of work - peak memory, bytes spilled, records in - against how many units shared the worker, then size from the observed peak plus room for bytes the engine never counted.

open as a page

Which properties of the branch above a step make recomputing it cheaper than keeping its result pinned in worker memory?

level: seniorimportance: should knowfreq 47%

basics

~20 s

Recomputation is cheap when every step in the branch is computed locally over a durable, re-readable source. It gets expensive when the branch crosses a step needing records from other workers, or reads a source that cannot be read twice identically.

open as a page

A long job pins several intermediates and its later steps spill heavily on unchanged input — what should you check first?

level: seniorimportance: should knowfreq 44%

basics

~20 s

Check whether earlier pins were ever released. A pinned result occupies retained-result memory until the job releases it or the engine drops it, so copies from finished phases can still be squeezing the operators of later steps on every worker.

open as a page

You double the number of pieces the input is cut into and the same single unit still runs out of room — what does that finding narrow the cause to?

level: seniorimportance: should knowfreq 52%

basics

~20 s

It rules out volume that the division rule can actually divide, and points at something the rule holds together: one group that must land in one place, or one record too large. The next lever is what the data is grouped by, not how finely it is cut.

open as a page

Why does handing records to a user-supplied function cost more than that function's own work, and what makes a cross-language one worse?

level: seniorimportance: should knowfreq 44%

basics

~20 s

Where the engine holds records packed, host-language code cannot read them: each record is decoded into objects going in and encoded back coming out. That cost is per record, so a trivial function can be dominated by it. A different language adds a process hop and a second conversion.

open as a page

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?

level: seniorimportance: should knowfreq 48%

basics

~20 s

Spilling 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.

open as a page

Two engines divide a worker's budget differently - one shared pool, one reserved block. What does each design buy and cost?

level: seniorimportance: should knowfreq 45%

basics

~20 s

A shared pool divided at runtime lets operator work and kept results take room the other is not using, at the price of one squeezing the other. A fixed reserved block is predictable but strands room it does not need.

open as a page

Your platform team's standard answer to a single unit running out of room is to raise every worker's memory — what does that policy cost?

level: principalimportance: should knowfreq 40%

basics

~20 s

It buys one fixed multiple of headroom and charges it on every worker process, in every job, for every run. It also trades parallelism when the machine total is fixed, and it hides the shape defect that will return at the next data size.

open as a page

A platform must standardise on one record representation for its shared pipelines - packed bytes or language objects - how do you decide?

level: principalimportance: should knowfreq 33%

basics

~20 s

Decide it as a capacity question with a measured answer: the footprint multiple sets how many machines every pipeline needs, while a packed default taxes each crossing into host-language code. Match the default to the dominant record shape, and leave a measured exception path.

open as a page

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%

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.

open as a page