When one worker process dies mid-flight, what work is redone under repeated small finite runs versus a runtime handling each record on arrival?
answer
- find the runtime's natural redo boundary
- a short run bounds its own redo
- continuous work has no natural boundary
- distance back to the last saved position
- arrivals keep coming during the redo
basics
~20 sThe execution shape fixes the unit of redo. Under repeated small finite runs, at most one interval's slice is re-executed. Under a runtime handling each record on arrival, the redo reaches back to the position the job last saved.
solid answer
~50 sA failure is rounded up to whatever boundary the runtime already has. **Repeated small finite runs** — where the runtime collects an interval's arrivals and executes an ordinary finite job over that slice — have a boundary every interval, so the worst case is re-executing one slice, and on many engines only the pieces the dead worker held. **Record-at-a-time processing** — each record handled as it lands, with context carried only in the job's own retained memory — has no natural boundary, so the runtime manufactures one by periodically saving its position and accumulated context; the redo is everything since that save. Engines differ on scope: some restart only the affected region of the job, others the whole running submission. In both shapes, arrivals accumulate during the redo, so the job must then sustain more than its input rate until it has caught up.
go deeper
Remember that a failure is rounded up to the nearest boundary the runtime has. A run per interval gives a boundary every interval; work that never finishes has none until one is manufactured.
Explain both bounds in mechanism terms — one slice against the span since the last save — and note that engines differ in whether they restart the affected region or the whole running submission.
Demonstrate the operational consequence: after the redo the job still owes the arrivals that accumulated, so recovery needs headroom above the input rate or the job never returns to its normal delay.
Treat the redo bound as a platform parameter with a price: shorter intervals and more frequent saves buy a smaller blast radius and cost fixed overhead, and the right point differs per pipeline.
## What the unit of failure means here When a machine is lost, the question is not only what was in flight on it but what the runtime has to go back to. Every runtime rounds a failure up to a boundary it already has, and the two execution shapes on the market have very different boundaries. - **Repeated small finite runs**: the runtime collects whatever arrived during a fixed interval and then executes an ordinary finite job over just that slice, over and over. A boundary arrives every interval, for free. - **Record-at-a-time processing**: each record is handled the moment it arrives, so nothing but the job's own retained memory carries context from one record to the next. There is no natural end, so the runtime has to manufacture a boundary by periodically saving its position in the input along with whatever it has accumulated. ## Repeated small finite runs: the interval bounds the redo - The largest thing that can be lost is one slice's worth of work, because everything before that slice was finished and recorded at an earlier boundary. - On many engines, only the pieces the dead worker was holding are re-executed, because the pieces that other workers already completed in the same run were recorded as they finished. - The magnitude of a failure therefore has the same order as the interval. A ten-second interval means roughly ten seconds of work to redo, plus the cost of starting the replacement work. - The unpleasant second-order effect is queueing: if the redo pushes a run past its interval, the next slice starts late, and the delay a reader sees grows until the job is back inside its budget. ## Handling each record on arrival: the last save bounds the redo - One submission runs and is expected never to finish, so there is no end to round up to. - The runtime periodically records where it had read up to and what it had accumulated. After a loss, work resumes from that record, so the redo spans everything that has arrived since. - **Engines differ in scope.** Some restart only the affected region of the job — the steps whose data flowed through the lost worker — while others restart the whole running submission. A design that works on one can be surprisingly expensive on the other. - A step that holds nothing between records is a different case again: on many engines the lost piece is simply recomputed from its inputs, because there was no accumulated context to restore. ## The comparison | | repeated small finite runs | handling each record on arrival | |---|---|---| | boundary that already exists | the end of every run | none; one is manufactured by saving periodically | | worst-case redo | one interval's slice | everything since the last save | | typical scope of the redo | the dead worker's pieces of that run | the affected region, or the whole submission, depending on the engine | | how to shrink the redo | shorten the interval, paying more fixed cost per run | save more often, paying the cost of each save | ## The obligation that follows any failure, in both shapes 1. Redo the lost work — a slice, or the span since the last save. 2. Work through everything that arrived while the redo was happening, which nobody stopped producing. 3. Sustain a rate above the input rate until that accumulation is gone; while it is not gone, the delay a reader sees is larger than the runtime's floor, and a job whose sustained rate is only equal to its input rate never recovers. This is why the unit of redo is a design parameter and not trivia: it sets how far behind a single machine loss puts you, and whether the job can climb back out without extra capacity. ## What this question deliberately does not cover How a saved position is taken consistently across many workers, what it contains, and what the outside world has already seen when work is repeated are all separate subjects. The point here is narrower and entirely about the execution shape: what boundary exists, and therefore how much work a failure costs. ## What an interviewer is listening for A candidate who names the boundary rather than the record. Weak answers say only that the lost record is retried, or assert one engine's recovery scope as though every engine of this class shared it. Strong answers connect the shape to the cost — a bounded slice against a span that grows with the interval between saves — and then remember the accumulation that has to be worked off afterwards.
- Why can one lost worker cost a continuous job far more than the work it was holding?Because the redo reaches back to the last saved position, not to the moment of the loss, and on some engines the restart covers more of the job than the failed part. Meanwhile arrivals accumulate, so the job must then run faster than its input rate to work off the gap.
- Does a shorter interval make failures cheaper as well as results fresher?It does bound the redo to a smaller slice, which is a real benefit. It also multiplies the fixed cost of starting and committing each run and the number of output commits, so the two effects are traded against each other rather than won together.
- What is redone when a step holds nothing between records?On many engines, just that piece of work, recomputed from its inputs — there was no accumulated context to restore, so the redo is proportional to the piece rather than to elapsed time. Engines that route work differently after a loss vary in how much surrounding work is repeated.
saying these in an interview costs you the question
- Says only the single record in flight has to be redone
- Assumes every engine restarts the whole submission after losing one worker
- Believes recovery always means restoring a saved position, whatever the job holds
- Forgets that records keep arriving throughout the redo
- Assumes a shorter interval only helps latency, never the cost of a failure
- Treats the redo as free because the cluster simply runs the piece again