skip to content

The Stage Barrier

The line across a job where no downstream worker may compute until every upstream one has finished - a stage barrier: fetching may overlap it, but the job is paced by its slowest producer.

on this pageshow

questions

4

In a finite job, why must a step that needs records from other workers wait for every producing piece to finish?

level: juniorimportance: must knowfreq 72%

answer

  1. a line, not a step
  2. completeness before compute
  3. the fold needs every producer's records
  4. slowest producing piece sets the pace

basics

~20 s

A step that needs records held by other workers cannot know it has them all until every piece of the producing step is done; computing sooner would publish a partial answer. That waiting line is a stage barrier.

solid answer

~50 s

Some steps need only what the worker already holds — test a predicate, drop a field — and those run straight through inside the piece. A step that folds or matches on a key needs records other workers are holding, so the job first performs a redistribution: every worker sends each record to whichever worker will handle that record's key. A consumer's key is complete only once *every* producing piece has emitted whatever it held for that key, so the job is cut by a **stage barrier**: a line across the job where no downstream worker may compute until every upstream piece of the producing step has finished. Two things draw that line — the operator genuinely needs completeness, and in the finite-job regime the exchange itself is often blocking. Fetching bytes may overlap producers that are still running; computing on them may not.

go deeper

for a junior

Recall the two shapes of step: one needs only the records already on the machine, the other needs records other workers hold. Only the second forces a wait, and the wait is about having all the input, not about speed.

for a middle

Explain the mechanics: the routing function makes equal keys meet, but completeness for a key is only known when every producing piece is done, and in the finite-job regime the exchange is often blocking on top of that.

for a senior

Show that you separate the two causes, because they have different remedies: an operator that needs completeness will always wait, whereas whether bytes move early is an exchange-mode property that differs between engines and changes only the transfer time.

for a principal

Frame it as the unit of scheduling risk on the platform: every line is a synchronisation point that converts one slow piece into whole-cluster idle time, so the count of lines in the standard job shapes belongs in your capacity and cost reasoning.

## Two shapes of step A running job is a graph of steps inside one program. (It is not the graph of scheduled jobs an orchestrator runs — same picture, different edges, different subject.) Within that graph the steps come in two shapes: - **Steps that need only what the worker already holds.** Testing a predicate, dropping a field, converting a value. The worker reads its own **piece** — one contiguous share of the input that a single worker processes on its own — and writes its own output. Several such steps in a row are usually run record by record inside the piece with nothing crossing the network. - **Steps that need records other workers are holding.** A total per customer, a match between two inputs on an identifier, an end-to-end ordered result. The records that must meet are scattered across the cluster, so they have to be moved first. That move is a **redistribution**: every worker sends each record it holds to whichever worker will handle that record's key, so the step can run at all. The first shape never makes anybody wait. The second is where the line appears. (Which operations fall into which shape is a neighbouring subject; here the point is only that the second shape exists.) ## Why there is a line at all Two distinct causes, and separating them is what a strong answer does. 1. **The operator needs completeness.** A count for a key is not *the* count until the last record bearing that key has been folded in. The routing function — the rule applied to each record's key that names its one destination worker — guarantees that equal keys meet, so a consumer knows *where* a key's records land. It cannot know *whether more are coming* until every producing piece has run out of input. Emitting earlier means emitting a number that is simply wrong. 2. **The exchange is often blocking.** In the finite-job regime many engines have each producer write its output into one bucket per destination on the local disks of the machine that produced it, and make that output collectable only once the producer has finished. Other engines push records across the network as they are produced. The first choice makes the wait explicit and makes re-running cheap; the second moves bytes earlier but does not remove the operator's need for completeness. Together they draw a **stage barrier**: a line across the job where no downstream worker may compute until every upstream piece of the producing step has finished. ## What the line does and does not mean - It is a **line, not a step**. Nothing executes "at" it; it is a constraint on when downstream compute may begin. - It is about **correctness, not a slow network**. A faster network shortens the transfer, not the wait. - It governs **compute, not transfer**. Bytes from producers that have already finished can be collected while others are still running. - It is a property of the **step boundary, not of a record**. A record that arrives early buys nothing by arriving early. - It is **not** the marker injected into the record flow so that every worker records its state at the same logical point. That marker is a recovery device; it exists to capture a consistent picture, not to hold compute behind a completeness requirement. ## Which steps sit behind a line | Step shape | Needs records other workers hold? | Waits for the producing step to finish? | |---|---|---| | Per-record projection or filter | No | No — runs through inside the piece | | Fold per key over the whole input | Yes | Yes — the fold is wrong until every producer is done | | Match two large inputs on a key | Yes | Yes — the side being matched against must be complete | | End-to-end ordered result | Yes | Yes — order is a property of the whole input | ## What the line costs - The step is **paced by its slowest producing piece**. The line lifts when the last piece finishes, not when the median one does, so wall clock follows the maximum and not the average. - Capacity behind the line is **idle, or merely collecting bytes**, which is how a cluster manages to look busy while producing almost nothing. - Where the regime materialises producer output, each line adds a **write and a read of the intermediate bytes** on top of the compute either side of it. ## Where the line is absent An endless input never finishes, so a job whose operators hand each record downstream as it is produced never reaches "every producer has finished" and therefore has no such line; its operators must emit from partial input using whatever bounding device the pipeline gives them. A large part of this market sits in between, running continuous work as a rapid succession of small finite jobs — and inside each of those small jobs the lines are back.

  • Does a chain of per-record steps with nothing moving between workers create a waiting line?
    No. If each step needs only the records the worker already holds, the whole chain runs inside the piece, usually record by record, and one worker finishing early does not hold anyone else up. The line appears only where a step needs records other workers are holding, so the records must be redistributed first.
  • Is it the producing piece or the whole producing step that must finish before a consumer computes?
    The whole step. A consumer's key can receive records from any producing piece, so completeness for that key is only known when every piece of the producing step has run out of input. One producer finishing tells the consumer nothing about whether more records for its keys are still coming.

Counting an election. A candidate's national total cannot be announced until every ballot box has reported, because the next box can change it — yet the boxes that have closed can already be driven to the counting hall while other polling stations are still open. Transport overlaps; the announcement does not.

saying these in an interview costs you the question

  • Thinks every step in a job waits for the one before it to finish.
  • Says the wait is the network being slow rather than a correctness requirement.
  • Believes a finite job can emit a key's total before all producers finish.
  • Confuses the waiting line with the marker used to snapshot a job for recovery.
  • Thinks adding workers to the waiting step removes the wait.
open as a page

While the producing step of a finite job is still running, what may the consuming side already do, and what must it not?

level: middleimportance: should knowfreq 48%

basics

~20 s

Collecting bytes from producers that have already finished may overlap producers that are still running; computing on them may not, because the consumer's input is not complete until the last producing piece is done. Transfer overlaps the wait, compute does not.

open as a page

Why does a continuous job that hands each record downstream as it is produced never reach a line where compute waits for all producers?

level: middleimportance: should knowfreq 52%

basics

~20 s

Because its producers never finish. A line across a job says compute waits until every upstream piece is done, and over an endless input that moment never arrives, so operators must emit from partial input instead of waiting for completeness.

open as a page

A finite job has three lines where compute waits for every producer; doubling the worker count changes its wall clock barely at all. Why?

level: seniorimportance: should knowfreq 50%

basics

~20 s

Each line lifts when its slowest producing piece finishes, not when the average one does. If every piece was already running at once, extra workers add capacity nobody was waiting for, so the job's wall clock stays roughly the sum of three maxima.

open as a page