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?
answer
- a rewind target, not a backup
- two parts, saved as one prefix
- input re-read from the recorded position
- the last 40 seconds recomputed
- output already written is not undone
basics
~10 sResuming 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 sA 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
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.
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.
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.
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