skip to content

Surviving Failure

What a job can lose when a machine dies and what it must never lose: recompute the piece, re-run the unit, or restore a saved picture - and what the outside world already saw.

on this pageshow

explore

questions

22

A long-running job crashes 40 seconds after its last recovery point. What does resuming from that saved picture restore, and what work is done again?

level: juniorimportance: must knowfreq 62%

answer

  1. a rewind target, not a backup
  2. two parts, saved as one prefix
  3. input re-read from the recorded position
  4. the last 40 seconds recomputed
  5. output already written is not undone

basics

~10 s

Resuming restores each worker's accumulated values and the recorded read position exactly as they stood at one coherent moment. Everything the job did in the following 40 seconds is read again and recomputed.

solid answer

~50 s

A recovery point is a durable copy of everything the job would otherwise lose — the running values each worker holds and the position it had reached in its input — written to storage outside any one machine while the job keeps running. On resuming, every worker loads its part of that copy and the input is re-read from the position recorded *inside it*, so the last 40 seconds of records go through the job a second time. Nothing after the recovery point survives: the picture is a rewind target, not a copy of the moment of the crash. This is the recovery model of a long-running job that accumulates values; a job holding nothing between records instead rebuilds a lost piece from its recorded derivation, and a finite job usually just re-runs the failed unit of work.

go deeper

for a junior

Recall the two halves of a recovery point — the values workers have accumulated and the position reached in the input — and that they are saved as one coherent moment rather than separately.

for a middle

Explain why resuming re-reads input: the position recorded inside the picture, not the point the crashed job actually reached, is what decides where work restarts.

for a senior

Separate the job's internal correctness from what the destination already saw. Replayed records produce their output again, and restoring a picture does not reach into a destination to undo it.

for a principal

Frame the 40 seconds as a cost the organisation chose. The real recovery time is that gap plus detection, obtaining machines, loading the parts and catching up, and that total is what anyone outside the team is promised.

## What a recovery point actually holds A **recovery point** is a durable copy of everything the job would otherwise lose, written while the job keeps running, so that a restart can resume from it instead of from the beginning of the input. The common industry word for it is a checkpoint, meaning exactly that mechanism: a copy taken of a live job, not an archive taken of a finished one. For a long-running job that accumulates values it holds two things, and they have to agree with each other: - **the values each worker holds** — running counts, partial aggregates, pending matches, timers — as they stood after some exact prefix of the input; - **the recorded read position** — how far into the input the job had durably got, so a restart knows where to resume. Both go to **durable shared storage for recovery points**: storage outside any one machine, reachable from whichever machine picks the work up. A copy on the local disk of the machine that died is not a recovery point, because the thing that failed is the thing holding it. *Consistent* here means one thing only: every saved part corresponds to the **same prefix** of the input. One worker's counts and another's include the same records and no others, and the recorded read position names the boundary of that prefix. ## What resuming does 1. The **coordinating process** — the single process holding the job's plan and tracking which units are done — finds the most recent capture that **completed**. A capture that was half written when the crash happened is discarded whole and an older complete one is used instead; a mixed picture is worse than a stale one. 2. Machines are obtained again, and each **worker** — one process on one machine running some of the job's units — loads the part of the copy assigned to it. 3. The input is re-read from the position recorded inside the picture, not from wherever the crashed job had actually got to. That live position was never durable, so it no longer exists. 4. Processing resumes, and the 40 seconds of records after the recovery point go through the job a second time. ## What survives and what is paid again | | Restored from the picture | Paid again after resuming | |---|---|---| | Accumulated values | as of the recovery point | the 40 seconds of updates to them | | Read position | as of the recovery point | the same input is read again | | Compute since the recovery point | nothing | all of it | | Output already written to a destination | not touched | may be produced a second time | The last row is the one candidates miss. Restoring makes the **job's own** values correct: each input record is reflected once in them. It does nothing about rows already written to a destination before the crash, which the replayed records will produce again. Internal correctness and what the outside world sees are two different claims, and the second one needs its own arrangement at the destination. ## This is one recovery model of three Asking how a distributed job recovers has no single answer, and the first thing to establish is what kind of job it is: - a job that **keeps nothing between records** recovers by rebuilding the lost piece from its recorded derivation — the engine's note of which inputs and which steps produced that piece — and needs no saved picture at all; - a **finite job** usually re-runs only the failed unit of work, and in the two-phase disk-to-disk model, where each phase materialises its output to durable files before the next phase reads them, that re-run starts from those files; - a **long-running job that accumulates values** is the case in the question: it holds something that cannot be recomputed cheaply, so it restores a picture and rewinds its input to match. Presenting the third as 'how recovery works' is the standard error, because it is false of the other two. ## Why 40 seconds is not the whole bill The distance back to the last recovery point sets how much input is reprocessed, but the time until the job is caught up is longer. The failure has to be noticed, machines have to be obtained, the saved parts have to be loaded, and only then does reprocessing start — while fresh input has continued to arrive throughout. The job then has to run faster than its input rate to close that gap, and how much faster decides whether catching up takes a minute or an hour.

  • If the crash happened while a recovery point was being written, what does the job resume from?
    The half-written one is unusable and is discarded whole, so the job resumes from the last capture that completed. A capture counts as complete only once every worker's part is durable and the coordinating process has written the metadata tying those parts together. A partial picture would mix values from different prefixes of the input, which is exactly the thing a consistent picture exists to prevent.
  • Why is a copy written to a worker's own local disk not a recovery point?
    Because recovery usually has to happen on a different machine. The failure that forces the restart is often the loss of that machine or its process, so a copy living only there is gone with it, and even where the machine survives the work may resume elsewhere. A recovery point has to live in storage reachable from whichever machine picks the work up, which is why writing one costs network and shared storage rather than just local disk.
  • How is this different from restoring a destination table to an earlier version?
    A saved picture of a running job is internal: it restores what the workers had accumulated and how far they had read, so the job can carry on. A stored table's committed version is a property of data at rest, restored by a mechanism the table format owns, and rolling one back does not rewind the job that wrote it. Both are casually called snapshots and they are not the same thing.

saying these in an interview costs you the question

  • Says a restart loses nothing, because the job writes recovery points
  • Treats the last recovery point as the job's state at the moment of the crash
  • Expects to reload the saved copy from the dead machine's local disk
  • Assumes restoring the picture also removes rows already written to the destination
  • Thinks every job needs a saved picture, even one holding nothing between records
open as a page

A job crashed while writing its results and restarted from its last recovery point — why can the destination hold some rows twice?

level: juniorimportance: must knowfreq 70%

basics

~10 s

Recovery is reprocessing: a restart re-reads input the job had already handled and writes it a second time. The destination remembers nothing of the first attempt, so a plain append adds those rows again.

open as a page

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%

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.

open as a page

A job crashes mid-run and restarts. Why does its recovery depend on the input being re-readable from a recorded position?

level: juniorimportance: must knowfreq 72%

basics

~20 s

Recovery is reprocessing: a restart re-reads input it already consumed. So the source must retain those records and let a reader resume from a durably recorded position — a source that forgets each record after handing it over makes the lost work unrecoverable.

open as a page

Parallel workers never pause at the same instant, so how does a running job save parts that add up to one coherent picture?

level: middleimportance: must knowfreq 58%

basics

~20 s

A marker injected at the sources travels with the records; each worker saves its own part as the marker reaches it and then passes it on, so every part covers the same prefix of the input.

open as a page

What makes a write key safe for a re-run to repeat, and which common key choices quietly break that?

level: middleimportance: must knowfreq 64%

basics

~20 s

A safe write key is computed only from fields the input record already carries, so a re-run derives the same key and the destination replaces the earlier row. Keys drawn from randomness, wall-clock time, attempt identity or a destination-assigned sequence all differ on the second run.

open as a page

Why must a job's recorded read position advance only after the output it covers is durable, and what does committing both as one unit buy?

level: middleimportance: must knowfreq 58%

basics

~20 s

Advancing the read position first turns a crash into permanent silent loss, because the restart resumes past output that was never written. Advancing it last turns the same crash into a duplicate. Committing output and position together removes the window where only one of them happened.

open as a page

A single unit of a batch job fails identically on every attempt while its neighbours succeed: what does that pattern tell you?

level: middleimportance: must knowfreq 62%

basics

~20 s

Identical failure on every attempt, across machines, points at a deterministic cause: the records that unit holds, or the code path they take. Retrying cannot fix it — only changing the input, the code or the unit's tolerance can.

open as a page

When a continuous job is asked to stop cleanly for a deploy, what must happen before the process exits?

level: middleimportance: must knowfreq 60%

basics

~20 s

Stop reading new input, let the records already inside reach the destination durably, advance the recorded read position to match, then write a restart point and exit. A hard kill loses no input but guarantees repeated work.

open as a page

Why is a job's recorded read position advanced only after the output it covers is durable, and what breaks if it is advanced first?

level: middleimportance: must knowfreq 64%

basics

~20 s

The ordering chooses which failure you get. Advancing the recorded read position first means a crash in the gap loses records nobody will ever reprocess. Writing output durably first means a crash re-does covered work — repetition, which the destination can be made to absorb.

open as a page

A job is being taken down for a deploy: what decides whether you can stop it and submit the new version, or must save first?

level: juniorimportance: should knowfreq 52%

basics

~20 s

Whether the job carries anything between records. A job that keeps nothing is stopped and resubmitted, provided its recorded read position is durable. A job holding accumulated results must write a restart point before it exits, or those results are gone.

open as a page

Why does a runtime that cuts an endless input into repeated small finite jobs need no marker flowing with the records?

level: middleimportance: should knowfreq 42%

basics

~20 s

Because the seam between two little jobs is already a moment with nothing in flight: all work for the completed piece has finished and none has started for the next, so the picture can simply be written there.

open as a page

A job's input is a directory of files rather than a positioned log. What must hold for a restart to resume it correctly?

level: middleimportance: should knowfreq 45%

basics

~20 s

Files must be immutable once visible, become visible in one atomic step, and survive until the job can no longer need them. The recorded position is then a set of file identities, and each identity must change whenever the content does.

open as a page

A worker fed by two upstream inputs receives the saving marker on one input long before the other. What are its options, and what does each cost?

level: seniorimportance: should knowfreq 38%

basics

~20 s

Two options: hold back the input that already delivered its marker until the other arrives, paying a pause that spreads upstream; or save at once and write the still-arriving records into the picture, paying size and a slower resume.

open as a page

What do you weigh when choosing how often a long-running job writes a recovery point, and what makes each one expensive?

level: seniorimportance: should knowfreq 50%

basics

~20 s

Weigh what one capture costs — the slowdown while workers save and the bytes pushed to shared storage — against the input a crash makes you reprocess, which is about half an interval on average.

open as a page

The destination is append-only, offers no transaction and enforces no key — how do you still make a replay leave no extra trace?

level: seniorimportance: should knowfreq 50%

basics

~20 s

Three moves remain: publish staged output in one visible act so abandoned attempts never appear, stamp every row with an identifier derived from the input record and collapse copies at read time, or record intent before an effect that cannot be retracted. The fourth honest answer is to accept duplicates and say so.

open as a page

A nightly job with a very high per-unit attempt cap now finishes four hours late but green: what is the cap hiding?

level: seniorimportance: should knowfreq 46%

basics

~20 s

A high attempt cap converts a broken job into a merely slow one. The extra four hours are retry time: some failure recurs, is paid for on every attempt, and never surfaces because the final status is green.

open as a page

A worker is lost five hours into a long run: when does that cost one re-run, and when must a whole region of the job rewind?

level: seniorimportance: should knowfreq 55%

basics

~20 s

Blast radius follows where the lost results lived. Durable outside the machine: only the in-flight unit re-runs. On its local disk or in its memory: the producing units re-run too. Held as accumulated state: a whole connected region rewinds to its last recovery point.

open as a page

How does a restart point an operator asks for before stopping a job differ from the ones written automatically?

level: seniorimportance: should knowfreq 44%

basics

~20 s

Different moment, different lifetime, different purpose. Automatic saves happen on a timer while the job runs and may be discarded when it stops; a restart point taken on purpose is written once after the drain, for a build not yet deployed.

open as a page

Your only input is a push feed that keeps no history. How do you make a job over it recoverable, and what does that cost?

level: seniorimportance: should knowfreq 52%

basics

~20 s

Put a durable, re-readable copy in front of the feed and have the job read that instead. Recoverability then starts at the copy, so the unrecoverable segment moves to the trivial hop that writes it — it shrinks, it does not vanish.

open as a page

Finance reports double-counted revenue from a pipeline the team believed could not repeat an effect — how do you decide which destinations to make repeat-proof and what the next notch costs?

level: principalimportance: should knowfreq 42%

basics

~20 s

The strong promise names an effect at one named destination, not a delivery across a pipeline, so it has to be decided hop by hop. Inventory each effect by whether it can be retracted and who already consumed it, then price the next notch in freshness, storage and coupling.

open as a page

You are budgeting the outage window for a planned stop and resume of a continuous job: what goes into it, and what must it fit inside?

level: principalimportance: should knowfreq 36%

basics

~20 s

Five spans: draining, writing the restart point, the deploy and getting machines, reading the saved copy back, and catching up the input that piled up. It must fit inside the source's retention and whatever freshness the consumers were promised.

open as a page