skip to content

A worker dies mid-piece in a 400-piece finite job: what does the engine re-run, and what must be true for that to be correct?

level: juniorimportance: must knowfreq 70%

answer

  1. smallest independently re-runnable thing
  2. one piece, not the whole job
  3. fresh attempt, any free machine
  4. partial output discarded, never merged
  5. only if input re-reads unchanged

basics

~20 s

The engine re-runs only the failed unit of work — the smallest piece it hands one worker — on whichever machine has capacity. That is correct only when the piece can be recomputed from re-readable input and its partial output is discarded.

solid answer

~50 s

The retry unit is the **unit of work**: the smallest thing the engine hands to one worker and can re-run independently, which in a finite job is usually one piece of the split input. When the worker holding it dies, the coordinating process — the single process that holds the job's plan and tracks which units are done — marks that unit failed and schedules a fresh attempt, typically on another machine, while the 399 pieces that already finished are untouched. Three things make that sound: the unit's input can be read again from the same position, its output depends only on that input rather than on anything the dead worker had accumulated, and whatever the failed attempt half-wrote is thrown away rather than mixed into the result. Where any of the three does not hold, the engine has to redo more than one unit.

go deeper

for a junior

Recall that the engine's retry unit is one piece of work, not the whole job, and that the failed attempt's half-written output is thrown away rather than completed or merged.

for a middle

Explain the three preconditions — input readable again from the same position, output determined only by that input, and partial output that can be abandoned — and say what goes wrong when each one is missing.

for a senior

Show that you have seen the coincidence break: name a job shape where losing one worker forces far more than one unit to run again, and say which number told you that was happening.

for a principal

Frame it as a price: the retry unit sets what one machine failure costs, and buying a smaller blast radius means materialising more intermediate results, which is paid for on every successful run too.

## Three different things called "the unit" A distributed job has at least three units, and retry granularity is about exactly one of them. - **The unit of parallelism** — one piece of the split input, which decides how much of the cluster works at once. - **The unit of retry** — the smallest thing the engine will re-run independently on another machine after a failure. That is the subject here. - **The unit of recovery** — the region of the job that has to rewind together to a *recovery point*: a durable copy of everything the job would otherwise lose, written while the job keeps running, so a restart can resume from it instead of from the beginning of the input. It is often much larger than either of the others. In a finite job that keeps nothing between records, the first two coincide: one piece of input, one unit of work, one re-run. Treating that coincidence as a law is the standard wrong answer, because it stops holding the moment the job accumulates state between records or keeps finished results in worker memory. ## What the engine actually does The **coordinating process** — the single process that holds the job's plan and tracks which units are done — stops hearing from the machine. (How it concludes the machine is gone is failure detection, a theory subject owned elsewhere.) It then: 1. marks every unit that machine was running as failed, and increments those units' attempt counters; 2. re-queues them as fresh attempts, placed wherever the cluster has capacity — **not necessarily on a different machine**, though several engines will steer away from a host that has failed repeatedly; 3. leaves alone every finished unit whose result is still readable from somewhere other than the dead machine. Step 3 is where designs diverge and where the whole blast-radius argument lives: an engine that materialises each step's output to storage outside any one machine has lost only the work in flight, while an engine that holds finished results in worker memory or on a local disk has also lost completed work and must re-run whatever produced it. ## The three preconditions for a per-unit re-run | precondition | what it means | what breaks without it | |---|---|---| | **Re-readable input** | the unit's slice of input can be read again, unchanged, from a position the engine recorded | there is nothing to re-run the unit against; what makes a source re-readable is a subject of its own | | **Output determined by that input** | the computation is a function of the records, not of a mutable external store or the current clock | the second attempt legitimately produces a different answer, so the retry silently changes the result instead of repairing it | | **Abandonable partial output** | whatever the failed attempt half-wrote can be thrown away | the surviving fragment plus a complete re-run means some records land twice | The third is usually arranged by **staged output promoted in a single step**: each attempt writes to a private location, and one final act makes the whole result visible at once, so an abandoned attempt leaves files nobody ever promotes. Where the destination offers no such staging — rows written one at a time into a live table, a call made to an outside system — the partial effect is already visible, and making a repeat harmless at the destination is a separate subject from re-running the unit. ## What the re-run costs - The unit's compute from the beginning: an attempt resumes nothing, because there is no saved partial progress inside a unit to restart from. - Re-reading the unit's input, which may mean pulling it across the network again. - Whatever warm local cache that attempt had built. - Wall clock for everything waiting behind it. If the step ends at a synchronisation point where every producing unit must finish before any consuming unit starts, one late re-run holds up the entire step. That last item is why a small retry unit is not automatically a cheap one. Cutting the input into more, smaller pieces makes each individual re-run cheaper, but a re-run that starts at the worst moment still costs a full round trip through everything queued behind it. ## What this is not - **Not the layer above re-running the whole job.** A scheduling system that runs a failed job again tomorrow morning is a different graph with its own rules, owned elsewhere; the engine's per-unit retry is exhausted entirely inside one run. - **Not the platform restarting the process.** Something that restarts a crashed process or replaces a machine hands the job a worker back; it restores none of the work, which the engine still has to re-run. - **Not a unit that is merely slow.** A worker that is behind has not failed, and the remedy for that has a different name and lives in a different part of this tree.

  • Does the engine guarantee that the re-run lands on a different machine from the one that failed?
    No. Most engines schedule the fresh attempt wherever capacity exists, which can include the machine that just failed; steering away from a host after repeated failures there is an extra policy some engines apply, not a property of retry itself. Placement is a scheduling decision and plays no part in the correctness argument, which is entirely about the input and the discarded partial output.
  • What breaks if a unit's computation reads a mutable external store or the current time?
    The re-run can compute a different answer from the first attempt, so a retry stops being a repair and becomes a silent change to the result — and two units retried at different moments can disagree with each other. Keep the unit a function of its input, or pin the external read for the whole run so any attempt at any time produces the same output.
  • Why does discarding the failed attempt's partial output usually cost nothing?
    Because each attempt writes to a private location and one final act promotes the whole result at once, so an abandoned attempt leaves files nobody promotes and a cleanup pass removes them. Where the destination has no such staging, the partial rows are already visible to readers, and the job needs a write that a repeat cannot double — a different subject from retry granularity.

saying these in an interview costs you the question

  • Thinks a failed unit resumes from where it stopped rather than starting over
  • Believes the re-run is always placed on a different machine than the one that failed
  • Assumes a half-written output from the failed attempt is harmless and left in place
  • Expects the engine to re-run the whole job whenever any single worker dies
  • Treats the piece of input a worker handles and the region that must rewind as the same thing