skip to content

Fifty minutes into a sixty-minute run the provider takes back one machine, and work that finished forty minutes ago is done again - why?

level: seniorimportance: should knowfreq 46%

answer

  1. completed does not mean consumed
  2. the product waited on local storage
  3. fetch fails, so the producer runs again
  4. more banked output later in the run
  5. the repair runs at risk too

basics

~20 s

Because that machine held intermediate output from earlier rounds that later pieces were still fetching. Losing it turns completed pieces back into work to be done, and the later the withdrawal, the more such output has accumulated.

solid answer

~50 s

A job of this shape runs in rounds: one wave of pieces finishes before the next begins, and the next wave's input is what the previous wave produced. On engines that materialise that output, it waits on the producing machine until consumers pull it. A piece is marked complete when it has written its share, not when the share has been read. Take the machine away and the consumers' fetches fail, so the steps that produced the missing shares are run again - and on engines that re-run the whole producing wave rather than a single piece, more finished work than that is discarded. The volume at risk grows as the run proceeds, which is why a late withdrawal costs a multiple of an early one. A runtime that pushes records with no materialisation loses nothing to fetch, and instead rewinds the affected part of the job to its last saved recovery point.

go deeper

for a junior

Recall that finished pieces are not automatically safe: their product may still be sitting on a machine, and if the machine goes the work is done again.

for a middle

Explain the window between a piece being marked complete and its product being fetched, and why that window is where a withdrawal turns finished work back into pending work.

for a senior

Reason about the growth: quantify what a machine had banked at the moment it went, note that the repair is on the critical path, and say that the repair itself runs on withdrawable machines.

for a principal

Argue about variance rather than averages. This capacity converts a predictable run time into a distribution with a tail, and the tail is what commitments are broken by.

## The job shape this question is about Not every distributed job is exposed here. A job whose pieces are independent - each reads its own slice of input, computes, writes its own result to durable storage outside the cluster - loses only what was in flight when a machine is withdrawn. The job in the question is the other kind. It runs in **rounds of work**: one wave of pieces must all finish before the next wave can start, because the next wave's input is what this wave produced. That dependency is the entire reason a late withdrawal is expensive. ## Completed is not consumed On engines that materialise intermediate output, a piece finishes by writing its product to storage local to the machine that produced it. The **coordinating process** - the single process that plans the pieces, hands them out and tracks what finished - marks the piece complete at that moment. Nobody has read it. Consumption happens later, when the next wave's pieces pull the shares they need across the network. So there is a window, and it can span most of the run, in which one piece of work is simultaneously: - reported as finished and counted in every progress display; - the only copy of data the job cannot proceed without; - resident on a machine the provider is free to take back. When the machine goes, the fetch fails. Nothing can conjure the bytes back, so the repair is to run again the steps that originally produced them - **lost-work recomputation**. The piece that finished forty minutes ago runs a second time, on a machine already busy with other work, and the wave waiting on it waits longer. ## Why lateness multiplies the cost Early in a run the machines hold almost nothing anyone else needs: little has been produced, and some of it has already been consumed. As the run proceeds, the volume of produced-and-still-needed output resident on each machine rises, and in a multi-round job the product of several rounds can be live at once. Two otherwise identical runs that lose a machine at different moments differ by exactly what that machine had banked. Three multipliers make the late case worse than the naive arithmetic suggests: 1. **Re-run width.** Some engines produce only the missing shares again; others discard and re-run the wave that contained them. The wider that unit, the more finished work is thrown away alongside the lost share. Which unit applies is the engine's own choice and a separate subject from this one. 2. **Serialisation.** The redone work is now on the critical path. The wave waiting for it cannot start, so the damage is elapsed time as well as machine-hours - and elapsed time is what a deadline is measured against. 3. **Exposure during the repair.** The recomputation runs on machines the job still holds, and those are withdrawable on the same terms. A busy period in the provider's capacity can produce a second withdrawal before the redone output has been consumed. That is what gives these runs a long tail rather than a fixed penalty. ## Where engines differ | how the runtime carries data between rounds | what a late withdrawal costs | what the repair looks like | |---|---|---| | materialises each round's output for the next to fetch | the produced output resident on that machine, which can be the product of several rounds | the producing steps run again, and the waiting wave is held up | | pushes records onward as they are produced | records in flight and the retained state of the lost workers | the affected part rewinds to its last saved recovery point and re-reads input from there | | runs continuous work as a rapid succession of small finite jobs | only what the small job currently in flight had produced | that small job's lost pieces are produced again, so the exposure resets frequently | The third row is worth noticing, because it shows that "late in the run" is a property of the execution model, not only of the clock: a runtime that keeps cutting the work into short finite units never accumulates an hour of unconsumed output in the first place. ## What shrinks the exposure - Landing durable output incrementally, outside the cluster, so finished work stops being at risk once it is written. - Keeping rounds short, so less produced output is unconsumed at any instant. - A **detached intermediate-output server** - a separate process that keeps serving a finished worker's output after that worker is gone - which covers a worker process exiting, but not the machine underneath it being withdrawn. - Avoiding a layout where a single machine holds a large share of what the next round needs, so one withdrawal cannot invalidate a whole wave. ## Three questions this is not Whether the engine retries one piece, the round or the job; what a restart resumes from; and a machine that is merely slow rather than gone. Each has its own answer, and conflating them with this one produces an argument that sounds complete and decides nothing.

  • Why can one missing share hold up a whole wave that was otherwise ready to run?
    A wave cannot start until the input it consumes exists. Consumers of the missing share wait, and on engines that re-run the producing wave rather than a single piece, everything in that wave waits with them. How wide the re-run unit is depends on the engine's retry granularity, which is a separate subject.
  • Does the recomputation itself run at risk?
    Yes. It runs on the machines the job still holds, which are on the same discounted terms, so a second withdrawal can arrive before the redone output has been consumed. That is why these runs show a long tail in completion time rather than a fixed, predictable penalty.
  • Does writing results durably as the job goes actually help?
    Where the job's shape allows it, yes. Work already written to durable storage outside the cluster is no longer at risk, so the exposure is only the work since the last durable write. A job that keeps everything inside the cluster until a final round exposes the whole run.

saying these in an interview costs you the question

  • Believes a piece that already reported completion can never be run again.
  • Thinks a withdrawal costs the same whenever in the run it lands.
  • Assumes survivors can read the lost output from the coordinating process's memory.
  • Says every engine keeps produced output somewhere that outlives the machine.
  • Treats the recomputation as free because the remaining machines were idle anyway.
  • Assumes a continuous job loses exactly what a finite job loses.