Parallel workers never pause at the same instant, so how does a running job save parts that add up to one coherent picture?
answer
- no usable global instant
- a cut that moves, not a time
- injected at the sources, flows in order
- save your part, then forward it
- usable only when every part is durable
basics
~20 sA 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.
solid answer
~40 sThere is no usable global instant: clocks disagree, and records in flight between workers would belong to no saved part. Instead a marker travelling with the records is used — a special element injected into the stream that each worker passes along after saving its own part. Because the marker keeps its place in the record order, a worker's saved part contains exactly the records ahead of it and none behind, so parts written by workers that never stopped together still add up to one coherent picture. The marker does not stop the job; it passes through, and work continues behind it. A runtime that serves an endless input as repeated small finite jobs needs no marker at all, because the end of each little job is already a moment with nothing in flight.
go deeper
Recall that a special element is injected into the stream and travels along with the records, and that each worker saves when it reaches it, so no clock is involved anywhere.
Explain the prefix property: the marker keeps its place in record order, so a saved part holds the records ahead of it and none behind, without any worker having to stop.
Show where this strains in production — one slow worker stretching the time until every part is durable, and captures that begin to overlap each other.
Weigh the two designs as a contract: a runtime with a natural quiet moment buys simplicity at a latency floor, while the marker buys low latency and brings reconciliation at fan-in with it.
## Why there is no usable simultaneous moment The obvious idea is to tell every worker to save what it holds right now. It fails for two independent reasons: - **Clocks disagree.** Machine clocks drift apart, so 'now' is not the same instant on two machines, and no amount of care makes it so. - **Records are in flight.** A record that has left one worker and not yet arrived at the next belongs to neither worker's saved part. On resuming it is either counted twice — re-read from the input while the upstream part already reflects it — or lost, because the sender treated it as sent and the receiver never saw it. The brute-force repair is to stop the whole job, let everything in flight land, then save. That is genuinely consistent and genuinely expensive: the cluster idles for the duration of every write. ## The marker travelling with the records A **marker travelling with the records** is a special element injected into the stream that each worker passes along after saving its own part, so that parts saved by workers that never stopped at the same instant still add up to one coherent picture. The mechanics: 1. The **coordinating process** — the one process holding the job's plan — asks the workers reading the input to begin a capture, identified by a number. 2. Each of those workers notes its **recorded read position**, saves it, and emits the marker into its output immediately after the last record it has already emitted. 3. A downstream **worker** that sees the marker on its input saves its own accumulated values to **durable shared storage for recovery points** — storage outside any one machine — then forwards the marker to its own outputs and carries straight on processing. 4. When every worker's part is durable, the coordinating process writes the metadata tying them together. Only then is the capture usable; until then it is a set of unrelated files. ## Why that is coherent The marker is not a moment in time; it is a **cut that moves through the dataflow**. Because it keeps its place in the record order on every edge it travels, each saved part contains exactly the records ahead of the cut and none behind it. Add the parts together and you have the job after some exact prefix of the input, which is all that consistency has ever meant here. No clock is consulted, and no worker needs to know what any other worker is doing. ## What the marker is not - **Not a stop.** Work behind the marker continues while a part is written, so throughput dips rather than falling to zero. That is the whole difference from the brute-force pause. - **Not the synchronisation point a wide step imposes**, which is owned elsewhere: that one makes every producing unit finish before any consuming unit starts, and it genuinely does hold the job up. The marker passes through workers; it does not gate them. - **Not a promise about the destination.** It makes the job's own values reflect each record once after resuming, and says nothing about output already sent. ## Where runtimes differ This mechanism belongs to **record-at-a-time processing** — a runtime where each record moves through the whole job as it arrives and workers hold values between records, so there is no natural quiet moment to save at. Other shapes in this class do not need it: | Runtime shape | How a coherent picture is obtained | |---|---| | Record-at-a-time | A marker travelling with the records | | Repeated small finite jobs | The seam at the end of each little job, where nothing is in flight | | Two-phase disk-to-disk | No running picture at all: each phase is materialised to durable files and a failed unit re-runs from them | So 'a marker is how it is done' is true of one part of this market and false of the rest, and an answer that does not say so is describing one product in neutral words. ## What goes wrong in practice - A worker fed by **several inputs** sees the marker arrive at different times on each, and has to reconcile them before its part is a clean prefix. - One slow worker stretches the time until every part is durable, so the capture takes longer even though nothing is broken. - A capture that never completes leaves the job resuming from an older one, which quietly widens how much input a restart reprocesses.
- Does the marker guarantee a destination sees each record's effect once?No. It makes the job's own accumulated values reflect each input record once after resuming, which is an internal claim about the job. Output sent to a destination before the crash is still there, and the replayed records will send it again. Turning that repeat into something invisible needs an arrangement at the destination, such as a write key computed only from the input record, and that is a separate mechanism from the capture.
- What happens to a capture if one worker never finishes saving its part?The capture never becomes usable. The parts are tied together by the coordinating process and the picture counts as complete only when all of them are durable, so one slow or stuck worker holds up the whole thing. Runtimes typically abandon a capture that exceeds a time limit and try again at the next marker, which is why a steadily rising capture duration is an early warning rather than a curiosity.
saying these in an interview costs you the question
- Says every worker simply saves at the same wall-clock time
- Thinks the marker stops the whole job while parts are written
- Ignores records in flight between workers when the picture is taken
- Believes a partly written picture can still be resumed from
- Assumes every runtime needs a marker, including one running repeated small finite jobs