skip to content

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%

answer

  1. capacity is not the constraint here
  2. memory is owned per process
  3. idle machines cannot lend memory
  4. compare the failing unit against the median

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.

solid answer

~40 s

A **worker** is one operating-system process, on one machine, running part of the job and owning a fixed amount of memory nothing else can borrow. A **unit of work** is one worker's share of one step, scheduled and retried on its own — and several of them usually run at once inside a single worker, dividing that one budget. So the failure is local. Idle machines elsewhere hold memory the failing process is not permitted to touch, which rules out cluster capacity as the cause and rules out a flaky host too, because the same unit failed on every attempt. What it leaves is a property of the data that reached that unit: more records than its neighbours got, or one thing inside it that has to be held whole.

go deeper

for a junior

Recall that a worker is one process with a fixed amount of memory nothing else can borrow, and that a job can die of a shortage on one machine while the cluster is mostly idle.

for a middle

Explain why the memory does not pool: the working set lives in one address space, so parallelism and per-process room are different resources. Then name the two data causes and the numbers that separate them.

for a senior

Show the diagnosis sequence out loud — failing unit against median unit, records in, bytes read, bytes spilled — and state which remedy you would try before touching the fleet's memory shape.

for a principal

The angle is policy: how often your platform answers this signature by raising everyone's memory, what that costs across all jobs, and whether per-job overrides with an expiry are the better default.

## The shape of the failure A **worker** here means one operating-system process, on one machine, that runs part of the job and owns a fixed amount of memory that nothing else can borrow. A **unit of work** is one worker's share of one step of the job — scheduled on its own, retried on its own, and usually several of them running at once inside a single worker process, dividing that one budget between them. **Running out of room** is worth splitting into the two failures it conflates: the engine refusing an allocation it was tracking, and the platform killing the whole process outright for crossing a ceiling the engine never saw. The pattern described — hundreds of units finishing quickly, one failing, retrying, failing again at the same point, and the rest of the cluster doing nothing — is the most diagnostic single picture in this whole subject, less for what it proves than for what it removes from the list. ## Why the idle machines are irrelevant Memory in a cluster of this kind is not pooled. Each worker process is given a budget when it starts, and the operators inside it borrow from that one budget while they run. A worker three racks away with 60 GB free cannot lend a byte to the process that is short, because the failing operator's working set — the records it is sorting, grouping, or holding while the step runs — lives in one address space on one machine. No system of this class serves one process's allocation out of another process's memory. Two consequences follow, and both are commonly got wrong in interviews: - **Adding machines does not help.** More machines add parallelism, not per-process room. If one unit of work needs more memory than one process has, a cluster of twice the size fails in exactly the same place at the same point. - **Cluster-level totals are the wrong number to quote.** "The cluster has 2 TB" is true and useless here. The number that matters is one worker's budget divided by how many units of work run concurrently inside that worker. ## What the pattern rules out, and what it leaves | Observation | Rules out | Leaves open | |---|---|---| | Hundreds of units finished | The step's logic being wrong for all data | Something specific to what reached one unit | | The same unit fails every attempt | A transient fault or one bad machine | A data-determined cause the retry reproduces | | The rest of the cluster is idle | Cluster capacity, and cluster-wide contention | One process's budget against one working set | | It fails at the same point each time | A race or a timing effect | A deterministic working set that does not fit | A retry that lands on a *different* machine and fails identically is not bad luck repeating itself; it is the strongest evidence available that the cause travelled with the data rather than with the host. ## The two causes that produce this picture 1. **The unit received far more than its share.** The step was fed by a **redistribution** — a step that cannot be computed from the records one worker already holds, so every worker writes its output split by destination and every worker fetches its share — and the rule choosing destinations sent an outsized fraction of the records to one of them. Why one grouping key carries such a share, and what to do about that key specifically, is a subject of its own and not this one. 2. **The unit holds one thing that cannot be divided at all.** A single group the operation has to hold whole, or one individual record that is enormous. Here the unit's total volume can be unremarkable and the failure still certain, because the indivisible thing alone is bigger than the room. Telling them apart is the next move, and the run's reported numbers do it: whatever the engine reports per unit of work after a run — records in, bytes read, peak memory, bytes written to local scratch disk — read for the failing unit and for a median unit beside it. Engines differ substantially in how much of this they report and under what names, but all of them report something per unit, and a memory claim is settled there rather than by arithmetic over the input size. ## What is actually expected of you here Not a fix. What the question is testing is whether the conversation goes to the right place: - Say that memory is owned per worker process, so an idle cluster is not evidence about capacity. - Say that four identical failures at the same unit make a flaky machine unlikely. - Ask for the failing unit's own numbers next to a normal unit's, rather than for the input's total size. - Know that raising every worker's memory is a real option and the expensive one: it is paid on every process in the cluster for the whole run, and it buys a fixed multiple of headroom for a problem that is usually a property of how the data was divided. That last point separates a useful answer from a reflex. Raising the per-worker budget may well make tonight's run pass. It just returns the same failure the moment the data grows past the multiple you bought.

  • The same unit fails on a different machine on each attempt. Does that change the diagnosis?
    It strengthens it. A fault belonging to one machine would be left behind by a retry elsewhere; a failure that follows the data onto fresh hardware is data-determined, which is exactly the case where more or healthier machines change nothing.
  • Several units of work run at once inside one worker. How does that change the arithmetic?
    They share the one process budget, so the room available to any single unit is roughly the budget divided by the concurrency, minus whatever the process uses outside the engine's accounting. Running fewer units per worker is therefore a genuine lever, at the cost of parallelism.

saying these in an interview costs you the question

  • Says the cluster needs more machines
  • Quotes total cluster memory as the relevant number
  • Treats four identical failures as flaky hardware
  • Assumes idle workers absorb the failing unit's overflow
  • Reaches for a memory setting before reading the failing unit's numbers