A finite job has three lines where compute waits for every producer; doubling the worker count changes its wall clock barely at all. Why?
answer
- sum of maxima, not total over workers
- extra capacity only absorbs queueing
- per-line costs ignore cluster size
- count the lines before tuning anything
basics
~20 sEach 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.
solid answer
~50 sThe job's wall clock is not the total work divided by the workers. It is, near enough, the sum over its lines of the slowest producing piece at each line — because a **stage barrier** is a line across the job where no downstream worker may compute until every upstream piece of the producing step has finished. Adding workers helps only while pieces are queueing for capacity. Once every piece of a step is already in flight, the step still ends with its longest piece, and the new workers sit behind the line collecting bytes or doing nothing. Two further costs are per line rather than per worker: the round of writing producer output and collecting it, and the synchronisation itself. So the levers are the number of lines and the length of the longest piece at each — not the size of the cluster. Why one piece is long is a separate diagnosis with its own owners.
code
python · 19 lines# Per-piece durations (seconds) for each step of a finite job.
# A line between steps: the next step computes only when the
# slowest piece of the previous one has finished.
steps = [
[30, 31, 29, 30], # even pieces
[12, 11, 600, 13], # one piece far longer than the rest
[20, 22, 21, 19], # even again
]
def step_duration(piece_seconds, workers):
if workers >= len(piece_seconds):
return max(piece_seconds) # every piece runs at once
return sum(piece_seconds) / workers # rough: pieces queue for capacity
def wall_clock(steps, workers):
return sum(step_duration(s, workers) for s in steps)
print(wall_clock(steps, 4)) # 653.0 -> 31 + 600 + 22
print(wall_clock(steps, 8)) # 653.0 -> unchanged; the pieces already fitgo deeper
Remember that a job's run time is set by the slowest piece at each waiting line, not by the average piece. That is why a cluster can be mostly idle and the job no faster.
Explain the arithmetic out loud: a sum of maxima plus a per-line round of writing and collecting output, and say precisely when extra capacity helps — only while pieces are still queueing.
Show judgment about interventions: count the lines first, establish whether the step is capacity-bound or maximum-bound, and hand off the cause of a long piece to the right diagnosis instead of buying machines.
Own the structural call: reducing the number of expensive movements, and reusing one when a pipeline needs it twice, changes the cost curve for every future run, whereas cluster size is a recurring bill that buys nothing behind a line.
## Where the wall clock actually goes For a finite job cut by lines where downstream compute waits for all producers, a usable model of run time is: > wall clock ≈ Σ over lines ( time of the slowest piece in that step + the round of writing and collecting that step's output ) Everything in that expression is a maximum or a fixed per-line cost. Neither responds to cluster size the way total work does. This is why a job can be running on a cluster that is 90% idle and still take exactly as long as it did yesterday. ## Why more workers stop helping Adding workers shortens a step only in one situation: when pieces were queueing because there was not enough capacity to run them all at once. Once each piece has somewhere to run, - the step still ends when its **longest** piece ends, and that piece's duration is a property of the piece and the machine, not of how many other machines exist; - capacity behind the line has nothing to compute, so the extra workers can only be collecting bytes, which is bounded by bandwidth rather than by worker count; - the per-line round of writing and collecting output does not shrink either — in fact more workers usually means more destinations, so the producing side writes more, smaller buckets. The common misreading is that the engine will notice the long piece and divide it. Some engines do revise the rest of a plan after measuring what a finished step actually produced, and some can re-cut a step's output for the step that follows — but none of that lifts a line for a piece that is still running, and in a continuous job that kind of re-planning largely does not exist. Treat it as a capability that varies, not as a property of the model. ## What each additional line costs 1. **A synchronisation point.** One slow piece is converted into idle time for every worker behind the line, so the cost of unevenness is multiplied by the number of lines it passes through. 2. **A round of intermediate bytes.** Where the regime materialises producer output, each line writes the step's output to local disk and reads it back across the network. That is real time and real device wear that no amount of compute capacity removes. 3. **A failure surface.** A line is a place where work already done is held, and where losing a producer's output means producing it again. Counting the lines in a job is therefore one of the cheapest diagnostics available: it tells you how many maxima you are paying for before you look at a single number. ## Which lever belongs to whom | Symptom | Lever | Whose subject it is | |---|---|---| | One piece far longer than the rest | Diagnose whether the data or the machine caused it | Unevenness and slow-worker diagnosis | | Too many or too few pieces per step | Choose the count against bytes and available threads | How the work is split | | A movement that should not exist | Arrange inputs so matching records already share a worker, fold before the move, or copy a small input everywhere | The movement-avoidance leaves of this category | | Three lines where one would do | Restructure the job so the expensive movement happens once | A design decision, informed by the count | The row that matters for this question is the last one: the count of lines is the thing this subject owns, and every other row is a hand-off. ## What to say in the interview State the model — wall clock is a sum of maxima plus per-line overhead, not total work over workers. Say what adding workers can and cannot buy: it can absorb queueing, it cannot shorten the longest piece, and it cannot remove a line. Then be explicit about the regime: this arithmetic is the finite-job story. A continuous pipeline that hands each record on as it is produced has no such lines, and its equivalent question is whether sustained throughput exceeds the input rate — a different measurement with different remedies. Finally, resist the reflex to blame the cluster. Doubling capacity is the most expensive intervention available and the one least likely to help a job whose shape is dominated by synchronisation. The cheap interventions are structural: fewer lines, and pieces that finish at roughly the same time.
- When would doubling the worker count actually cut the wall clock?When pieces are queueing for capacity — many more pieces than places to run them. Then the step is capacity-bound and extra workers run more pieces at once. The moment every piece of a step is in flight, the step's duration is its longest piece and further capacity changes nothing behind that line.
- Does removing one of the three lines save only the time spent waiting?No. It also removes that line's round of writing producer output and collecting it across the network, and one synchronisation point at which one slow piece idles the whole cluster. Where the regime materialises intermediate output, those two costs are often larger than the wait itself.
saying these in an interview costs you the question
- Computes wall clock as total work divided by worker count.
- Says more machines always help a job that is running slowly.
- Assumes the engine will detect a long piece and split it mid-flight.
- Counts only the compute and ignores the per-line write and collect.
- Applies the sum-of-maxima model to a record-at-a-time pipeline.