skip to content

A machine holding a finished producing step's written buckets is lost while later steps run — what does recovering that cost?

level: seniorimportance: should knowfreq 54%

answer

  1. intermediate output, not the job's output
  2. local disks, usually a single copy
  3. the recipe survives, the bytes do not
  4. an exited worker is not a lost machine

basics

~20 s

The producing work again, not just a transfer. Those buckets are intermediate bytes on one machine's local disks, usually in a single copy, so nothing can be re-read: engines that kept a description of how the piece was made re-run it.

solid answer

~50 s

A redistribution's intermediate output is not the job's output. Final results go to the shared store — durable storage every worker can reach — and are usually replicated there; the per-destination buckets a producer writes sit on that one machine's local disks, typically in a single copy, because replicating every intermediate byte often costs more than making it again. So when the machine goes, the bytes go. What survives is the *recipe*: the engine kept a description of how that piece was produced, and can re-run the producing work from that step's own inputs. The bill is therefore recomputation, not a retry of the transfer — and if the lost step's input was itself the output of an earlier movement, the same question is asked one level further up. Engines differ sharply here: some recompute automatically, some replicate or offload intermediate output, and a runtime that pushes records as it makes them has no written share to lose at all.

go deeper

for a junior

Recall that the buckets a producing step writes are temporary, live on one machine's local disks, and are not the durable copy that the job's final output gets.

for a middle

Explain why the loss costs the producing computation rather than a repeated transfer: each producer's share is unique, there is no second copy, and what the engine kept is a description of how to make it again.

for a senior

Separate the worker's lifetime from the machine's, and say what each arrangement survives. Then note what varies: automatic recomputation, offloaded intermediate output, and hand-offs that store nothing recover in three different ways.

for a principal

The trade is platform-wide: paying to replicate or offload intermediate output on every movement against occasionally paying recomputation, and that arithmetic shifts sharply once your machines can be reclaimed mid-job.

## Two kinds of output, and only one of them is precious A job produces two very different classes of bytes. - **Final output** goes to the **shared store** — durable storage every worker can reach — and is normally replicated there, because losing it means losing the answer. - **Intermediate output** is what a redistribution leaves behind: one bucket per destination, written by a producing unit to the local disks of the machine it ran on. It exists only so the collecting step can pick up its share, and it is normally held in a single copy. That single copy is a deliberate trade, not an oversight. Replicating every intermediate byte would add a full extra network crossing and extra storage to every movement in every job, to insure against a machine loss that usually does not happen. Most designs decide that making the bytes again on the rare occasion is cheaper than copying them always. ## Why the bill lands on the producing work When the machine holding a share disappears, a collector asking for that share gets nothing, and there is no second copy to ask for instead. Each producer's bucket is *unique*: it holds the records that this producer made for this destination, and no other producer has them. What the engine still has is the description of how that piece was produced. So the repair is to **run the producing work again** — read that step's own input, redo its computation, re-route its records, re-hold its buckets — and only then can the waiting collector be served. Three consequences follow: 1. The expensive part is the **recomputation**, not the transfer. A transfer is seconds; the producing work can be minutes of processor time over a full piece of input. 2. If the lost step's input was itself the output of an earlier movement, the earlier shares may be needed again too, and the question repeats one level up. A loss late in a deep job can therefore unwind a long way. 3. The **later a loss happens, the more it costs**, because more producing work is resting on those disks by then. ## Surviving the worker is not the same as surviving the machine Two distinct lifetimes are easy to conflate: the worker process that produced the bytes, and the machine whose disks hold them. | arrangement | survives the producing worker exiting | survives the machine going away | what it costs | |---|---|---|---| | the producing worker serves its own output | no — the worker must stay alive | no | idle producing capacity cannot be released | | a separate process on the machine keeps serving the output | yes | no | one more process per machine to run and operate | | the output is written somewhere remote or replicated | yes | yes | every intermediate byte crosses the network at least once more | | nothing is stored; records are pushed as produced | there is no stored share at all | there is no stored share at all | recovery is an entirely different mechanism | The second row is the one candidates most often get wrong in both directions. A process that outlives the producing worker is genuinely valuable — it lets a job hand back workers it no longer needs while their output is still being collected — but it does nothing whatsoever about a lost machine. The disks are still that machine's disks. ## What varies between engines, and how to say it This is a place where a confident sentence about how it works is usually a sentence about one product: - whether intermediate output is replicated or offloaded at all, and whether that is a choice the platform makes or the job does; - whether a separate serving process exists, so capacity can be reclaimed mid-job; - whether losing a share triggers automatic recomputation or simply fails the job; - and, in a runtime that pushes each record downstream as it is made, there is **no written bucket to lose**. What a machine loss takes there is the records in flight and whatever the consumer had accumulated, and the job returns to a recorded point rather than re-collecting a share — a different recovery subject entirely. The safe formulation in an interview is to state the mechanism and then name the variation: *the share is on one machine's local disks, so its loss costs the producing work again, on engines that can recompute it from what they remember — and there is no share at all where the hand-off stores nothing.* ## Where this stops This explains **why** a loss is expensive and **what** gets redone in principle. Exactly which unit is retried, how many attempts it gets, and where a retry resumes from are questions about failure granularity and belong with recovery, not with the mechanics of the exchange.

  • What does a separate process that keeps serving a producer's output after the worker exits actually buy you?
    It decouples the output's availability from the producing worker's lifetime, so a job can release workers it no longer needs while their shares are still being collected — which matters most on long jobs with wide early steps. It does not decouple the output from the *machine*: if the host goes away, the shares go with it.
  • Why not simply replicate intermediate output the way the shared store replicates final output?
    Because you would pay for it on every movement in every job, whether or not anything is ever lost. An extra copy means another network crossing and more storage for bytes whose whole purpose is to be consumed once and discarded. Some platforms do offer it, precisely for workloads on machines that are reclaimed often.
  • Why is losing a machine late in a long job worse than losing one early?
    Because more completed producing work is resting on local disks by then, and any of it that is still needed must be made again. If the lost step's inputs were themselves produced by an earlier movement whose shares are also gone, the recomputation walks further back up the chain.

saying these in an interview costs you the question

  • Assumes intermediate buckets are replicated like output in the shared store
  • Thinks a worker exiting always destroys the shares it wrote
  • Says the collector can just read the same share from another producer
  • Believes only the interrupted transfer has to be repeated
  • Treats every engine as materialising, ignoring hand-offs that store nothing
  • Thinks re-reading the source is the whole cost, not the computation