skip to content

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%

answer

  1. one process, one ceiling
  2. shares of the same fixed budget
  3. the copy sits where operators work
  4. less room means writing runs to disk
  5. spill figures rise on later steps

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.

solid answer

~50 s

A worker is one process with a fixed budget that several units of work share, and that budget is divided before any record arrives. A pinned copy is held in **retained-result memory**; the sorts, grouping tables and held join sides of the next step draw on **operator working memory**. They are shares of the same fixed number, so a large pin leaves the next operator less room, and an operator that cannot hold its working set starts spilling — writing part of it out, in sorted runs or per-key chunks, to disk attached to the worker and reading it back to finish the step. You have swapped one walk of the branch for extra disk traffic on every step that follows. Engines draw the region line differently: one splits a single pool dynamically between the two, one reserves a block it manages itself and leaves the rest to the language runtime, one barely has held results at all.

go deeper

for a junior

Remember that a worker is one process with a fixed amount of memory, and that anything kept for reuse is held inside that same amount rather than somewhere extra.

for a middle

Explain the chain: the copy occupies the retained-result share, the next operator gets less working memory, it writes part of its working set to local disk, and every following step pays that cost.

for a senior

Demonstrate that you judge a pin against what it displaces rather than against nothing — comparing per-unit spill and peak-memory figures before and after, and noticing when the engine is holding the copy and recomputing the branch anyway.

for a principal

Decide the platform default: whether pinning is available to authors, what worker shape everyone inherits, and how many units of work share one budget — because the same pin is harmless on one shape and fatal on another.

## One process, one budget 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. Several **units of work** — each one worker's share of one step, scheduled and retried on its own — run inside it at the same time and share that same memory. So whenever you state a budget, state in the same breath how many units of work are drawing on it: a worker with eight units of work running has eight concurrent operators dividing whatever operator space is left. The cluster's total size is beside the point. Memory is owned per process, so a cluster with terabytes free can still contain one worker process with nowhere to put a sort. That per-process budget is divided before the first record arrives. Two of the regions matter here: - **operator working memory** — what operators borrow while they run (sorting, building a grouping table, building the held side of a join) and give back when the step ends; - **retained-result memory** — what holds results the job was told to keep for later reuse. A pinned copy lands in the second. Whatever it occupies is not available to the first. That is the whole mechanism, and it is why a pin is a two-sided decision rather than a free optimisation. ## The displacement chain 1. You pin a large intermediate. Each worker keeps its own share of it. 2. The next step sorts or aggregates. Its operators ask for working memory and are given less than they had before the pin. 3. An operator that cannot hold its working set **spills**: it writes part of that set out — in sorted runs, or per-key chunks — to **local scratch disk**, disk attached to the worker machine whose contents nobody may read once the job ends, and reads it back to finish the step. 4. Spilling is planned degradation, not a defect. But it is not free: bytes written, bytes read back, extra merge passes, and in the worst case the same data spilled and refilled repeatedly. 5. Net effect: you saved `n - 1` walks of one branch and paid for it with disk traffic on every step that follows. | what the pin gives | what the pin takes | |---|---| | one walk of the branch instead of `n` | operator working memory on every worker, for as long as it is held | | stable read cost for the later readers | spill writes and reads on the steps after it | | a fixed cost you can measure once | a cost that scales with how long the pin is held and how many units of work share the worker | Below some size the trade is a clear win. Above it, the job is slower than the recomputation the pin was meant to avoid — on unchanged input, which is exactly the situation in this question. ## Engines divide the budget differently This is the part of the subject where the three lineages in this family disagree most, so no single sentence about "how much the pin took" is true everywhere: - one splits a single pool dynamically between operator work and results it was told to keep, so a pin can borrow into operator space and be pushed back out of it under pressure; - one reserves a block it manages itself as packed bytes and leaves the remainder to the host language runtime, so the interaction is between that managed block and the runtime's own allocation and reclamation pressure; - one has almost no notion of kept results at all and spills from a fixed sort buffer, so the displacement question barely arises. The safe generalisation is the one that holds across all three: the copy and the operators are inside one process with one ceiling. ## Reading it instead of guessing Settle this from **the run's reported numbers** — whatever the engine reports per unit of work after a run: - **spill figures** on the steps after the pin, before and after introducing it. Note which figure you are quoting: the size the data occupied in memory, or the bytes actually written. The two differ by however the engine encodes the data, often by a large multiple, and that difference is a fact about encoding rather than a bug. - **peak memory per unit of work**, which tells you how close the operators were running to their region boundary. - whether the **branch's own steps reappear** in the run. If they do, the engine dropped the pinned result under pressure and you are now paying for both the copy and the recompute. - **reclamation pauses** — the runtime stopping the worker's own work while it finds and frees objects nothing refers to — which lengthen as held objects accumulate on engines that keep results as host-runtime objects. ## What to do about it - **Pin less of it.** Drop the columns no reader uses before pinning; the copy shrinks in proportion. - **Pin a denser form.** Where the engine can keep a result as a packed byte layout it interprets itself rather than as ordinary language objects, the copy is several times smaller at the cost of a decode step per read. - **Let it overflow rather than displace.** Where the engine can keep the surplus on local scratch disk, a disk read per use is often cheaper than pushing every later operator into spilling. - **Release it** as soon as the last reader is done, so the displacement lasts phases rather than the whole run. - **Do not pin at all** when the branch is cheap to walk again — a filtered scan of a durable source often is.

  • The pin is held and the following steps also recompute the branch. How is that possible?
    On engines where a pin is a request rather than a guarantee, the engine can drop a pinned result to make room for operator work and recompute it on the next read. You then pay twice: the copy occupied memory long enough to push operators into spilling, and the branch is walked again anyway. It shows up as the branch's steps reappearing in the run's reported numbers while held bytes are still counted.
  • Does adding more workers relieve the pressure a pin creates?
    Only indirectly. The copy is distributed, so more workers means each holds a smaller share — but each worker also holds a smaller share of the input, and the number of units of work per process often stays the same. If the job is failing because one group must be held whole, more workers changes nothing. Dividing the data differently is usually the better lever than adding capacity.
  • Why can a worker be killed outright after a pin, at a number the engine says it never reached?
    The engine accounts for the regions it manages, but the process also uses memory it never counted — the language runtime, native and network buffers, and whatever user-supplied functions allocate. That off-budget overhead counts against the ceiling the platform enforces on the whole process, and the platform enforces it by killing the process rather than by failing one allocation. A pinned copy held outside the accounted regions widens that gap.

A desk of fixed size. Leaving the finished, tidied stack of results on the desk because two colleagues will want it saves them re-sorting the whole pile — but it leaves less surface for the sorting you are doing now, so you start shuttling working piles to the floor and back. Past a certain stack size, the shuttling costs more than the re-sort you saved.

saying these in an interview costs you the question

  • Retained results and operator work draw on separate, independent pools
  • The cluster has plenty of free memory, so the pin cannot be the cause
  • Spilling proves the job is broken and needs more memory everywhere
  • A pin costs nothing once the result has been computed
  • Raising the worker's memory is the only fix for a displacing pin