skip to content

One process finishes or dies whole, but a cluster run can end with some pieces done and others lost. What does that possibility cost a team?

level: seniorimportance: should knowfreq 55%

answer

  1. all-or-nothing is a single-machine luxury
  2. half the output may already be visible
  3. degraded is neither success nor failure
  4. reruns are safe only if writes repeat safely
  5. done is a fact about the destination

basics

~20 s

Partial failure: some work is durable and some died with its machine. The team now owns a write path and a rerun that stay correct when the same records arrive twice, monitoring that can tell degraded from failed, and a definition of done that is not an exit code.

solid answer

~60 s

On one machine the run either completes or it does not, and the failure is total and visible. A cluster adds states that have no single-machine equivalent: a worker vanished while the rest kept going, half the output is already durable at the destination, the run is degraded but still progressing, or the coordinating process is alive while much of its capacity is gone. Each of those turns into engineering work somebody has to fund. The write path has to tolerate the same records arriving a second time, because a rerun will produce them. Monitoring has to distinguish slow from dying, because the two look alike. Cleanup has to have an owner, because a half-written destination will otherwise be read as a whole one. How the engine itself recovers varies — a finite run usually re-executes the lost unit, a long-lived run restores from a recovery point, and an engine holding nothing recomputes the lost piece from its inputs — but the obligation at the destination is the same in all three cases.

go deeper

for a junior

Recall that a distributed run can end with part of its work finished and part lost, which one process on one machine simply cannot do, and that this is why output may be half-written.

for a middle

Explain the consequences: a write path that must tolerate repeats, a rerun whose safety depends on the destination, and a definition of done that is a check against the output rather than an exit code.

for a senior

Show that you have operated this. Name the degraded-but-running state, say how you alert on it separately from failure, who cleans up partial output, and how much progress is at risk given how the engine in front of you recovers.

for a principal

Own it as a standard: what every pipeline's write path must guarantee before it is allowed into production, who is on call for half-finished runs, and what the platform provides so each team does not solve repeat-safe writing on its own.

## A failure mode one machine cannot have Run a computation as one process on one machine and failure is simple: the process exits, and whatever it had not yet written is gone. There is one thing to watch, one exit code, one log. The failure is total, which is unpleasant but unambiguous. Spread the same computation over a **cluster execution engine** — a system you hand a whole program to, which splits its work across many machines and puts the results back together — and a new state becomes reachable, one no single process can enter: **part of the work is finished and durable while another part died with its machine.** This is not the cost of the failure itself. It is the cost of the *possibility*, and it is paid on every run, including the ones that succeed, in the form of work the team must do so that the bad case stays survivable. ## Where the cost actually lands 1. **At the destination.** If output is written as the run proceeds, a failure leaves the destination in a state no one designed: some records present, some absent, possibly some counted twice after a retry. The write path now has to be **idempotent** — safe to repeat — or it has to publish only at a point where everything is present at once. Either is real design work that the single-machine version never needed. 2. **In monitoring.** A worker that has vanished and a worker that is merely slow look the same from outside for as long as the timeout allows. A run can be *degraded* — fewer workers, still progressing — which is neither success nor failure, and something has to decide how long that is acceptable. 3. **In reruns.** "Just run it again" is a safe sentence on one machine and a loaded one here, because the second run may re-emit what the first already delivered. Whether a rerun is safe is a property of the destination, not of the engine. 4. **In the definition of done.** An exit code no longer settles it. Done becomes a statement about the destination — this many records, this marker present, this partition complete — and someone has to write that check. 5. **In the blast radius of one machine.** A machine reclaimed late in a long run can force work that already succeeded to be done again, so the cost of one loss is not proportional to how much of the run it held. | | one process on one machine | a distributed run | |---|---|---| | possible end states | done, or dead | done, dead, or partly both | | what failure destroys | everything not yet written | whatever was on the lost machine | | is a rerun safe? | usually, trivially | only if the destination tolerates repeats | | what "done" means | the process exited cleanly | the destination holds the whole answer | | who notices | the caller | monitoring the team had to build | ## What varies between engines The engine's own response to a lost worker is one of the places where this family disagrees most, so do not state one lineage's behaviour as the model: - A **finite run** with no long-lived memory usually re-executes the lost unit of work and, if its inputs are gone too, recomputes them from further back. - A **long-lived run that retains state** restores from a saved recovery point and resumes from there, which means the recovery point's age is the amount of progress at risk. - An engine that **retains nothing between records** can simply recompute the lost piece from its input, provided the input can still be re-read. What does *not* vary is the obligation at the edge of the system. Whichever of those three the engine does, the destination may see the same records more than once, and no engine can make an external write un-happen on your behalf. ## The honest interview answer Say the state exists, say why it cannot exist on one machine, and then price it: an idempotent or all-at-once write path, alerting that separates degraded from dead, an owner for partial output, and a completeness check that replaces the exit code. Add the organisational half — somebody must be reachable when a run is half-done at three in the morning, which is a rota the single-machine version did not need. A candidate who answers only "the engine retries it" has described the engine's job and missed the team's.

  • Why is a rerun a design decision here rather than an obvious remedy?
    Because the first attempt may already have delivered part of its output to somewhere outside the run's control. A rerun then re-delivers those records, and whether that is harmless depends entirely on the destination: a keyed overwrite absorbs it, an append or an outgoing notification does not. The engine cannot decide this for you.
  • Does a run that survives losing a worker mean no cost was paid?
    No. Surviving means work was redone, capacity was down while it was redone, and the run finished later than planned — which for a deadline-bound pipeline is its own kind of failure. It also means monitoring had to distinguish that slowdown from a genuine stall, which is the standing cost the possibility imposes on every run.
  • How should completeness be checked if an exit code no longer settles it?
    By asserting something about the destination rather than the run: an expected record or partition count, a marker written only after everything else, or a reconciliation against the input's own counts. The check belongs to the pipeline, not to the engine, and it should run before anything downstream is allowed to read.

saying these in an interview costs you the question

  • Says a failed run leaves the destination untouched
  • Treats retrying as free and always safe
  • Believes the engine guarantees an external write happens once
  • Cannot distinguish a degraded run from a failed one
  • Uses the exit code as the definition of done
  • Assumes every engine recovers a lost worker the same way