After fanning a scoring pipeline out across workers, why do results arrive out of order, and how do you restore it?
answer
- completion order replaces source order
- handed out in order, finish out of order
- carry the position through the fan-out
- the bound does not bound held results
- cheapest fix: an order-independent sink
basics
~20 sResults converge in completion order, not source order, because the work for each record takes a different amount of time. Restore order by carrying each record's position through the fan-out and re-emitting by position, or by making the write-back address records by key.
solid answer
~50 sOnce several records are being worked on at once, the merge or rejoin passes on whichever finishes first, so source order is gone by construction. There are three honest responses. Carry each record's **position** through the fan-out and add a reordering step that holds finished results until the next expected position arrives — cheap in code, but it needs an explicit window and a policy for a straggler. Make the sink **order-independent**, writing each result against its own key so position never mattered; this is usually the cheapest fix and the one to reach for first. Or give up the overlap for the records that must stay ordered. What you should not do is buffer the entire result set to sort at the end, which reintroduces exactly the memory cost a streaming batch exists to avoid.
code
pseudocode · 17 lines// tag before the fan-out, reorder after the convergence point
indexed = source(records).withIndex() // (0,r0), (1,r1), (2,r2), ...
scored = indexed.fanOut(
item -> innerOperation(() -> pair(item.index, score(item.record)))
.executeOn(workerPool),
maxInFlight = 8)
nextWanted = 0
held = emptyMap()
for each result in scored:
held.put(result.index, result.value)
while held.contains(nextWanted):
emit(held.remove(nextWanted))
nextWanted = nextWanted + 1
if size(held) > window:
fail("reorder window exceeded at position " + nextWanted)go deeper
Hold on to the cause: once several records are worked on at once, results come back in the order they finish, which is not the order they went in. Any ordering you need afterwards has to be rebuilt deliberately.
Explain the mechanism and the standard fix: attach each record's position before the fan-out, and reorder after the results converge by holding out-of-order ones until the next expected position arrives.
Demonstrate the production judgment: ask first whether the sink needs order at all, bound the reordering step with an explicit window, choose between stalling and failing on a straggler, and say how a failed parallel run is resumed.
Treat ordering as a contract with the consumer of the output. Decide whether downstream jobs may depend on positional order at all, since removing that dependency is usually cheaper and more durable than paying for reordering in every pipeline.
## Why order disappears A fan-out starts work for several records at once; a track split hands shares of the stream to different workers. In both shapes the convergence point — the merge, or the rejoin — emits **whatever finished first**. Per-record durations differ, for reasons you do not control: a larger record, a cache miss, a slower dependency, a worker descheduled mid-run. So the output sequence is in completion order, and completion order is not reproducible between runs. Two confusions are worth clearing up immediately: - **Elements are handed out in order.** What is undefined is the order they *complete* in, not the order they start. - **Order within one track is preserved**, because a track is itself a serial sequence. It is the rejoin across tracks that loses it. ## The three responses | Response | What it costs | When it fits | |---|---|---| | Carry the position and reorder at the rejoin | A held-results buffer plus a straggler policy | The sink genuinely needs positional order | | Make the sink order-independent | A key on every record; often none at all | Results are addressed individually | | Keep the ordered records serial | The overlap you were buying | A small ordered tail among independent work | **Check the second row first.** Much of the time the nightly batch writes each score against its record's identifier, and nothing downstream depends on the order the writes happened in. Then reordering is work you do not need. The requirement to preserve order is real only when the *position* of the output carries meaning — a positional file the next job reads, an append-only log consumed by index, a report whose rows are the input rows. ## Reordering by position, done properly Tag each element with its index before the fan-out, keep the index attached to the result, and place a reordering step after the convergence point. It holds out-of-order results in a map and releases them as the next expected index arrives. The naive version of that step has a real defect, and it is what a senior interviewer is listening for: 1. The fan-out's concurrency bound caps how many records are *being worked on*. It does not cap how many finished results are *waiting*. 2. When one record is slow, its slot stays occupied, but every other slot keeps completing and starting new records. 3. Those completed results flow to the reordering step and pile up, all blocked behind the single missing index. 4. Held results therefore grow with the input, not with the bound — the very unboundedness the bound was added to prevent, moved one stage downstream. So a reordering step needs its own **window**: a maximum distance the pipeline may run ahead of the oldest unfinished position, and a decision about what happens when that distance is exceeded. The realistic options are to stop signalling demand upstream until the straggler lands — trading throughput for a bounded footprint — or to fail the run with a diagnosable error naming the position that stalled. Silently growing is not an option, and neither is quietly dropping. ## What not to do - **Collect everything and sort at the end.** It works, and it costs the whole result set in memory — which is what a streaming batch exists to avoid, and it will pass every test that uses a small fixture. - **Sort by a value in the record instead of its position.** Records need not carry anything that reproduces input order, and a second full pass over the output costs another read and another write. - **Emit only when the oldest element completes.** That removes the overlap entirely and returns the pipeline to serial execution wearing a fan-out's clothes. ## Failure is where ordering hurts most A serial batch that fails has written a prefix of the results; you know where to resume. A parallel batch that fails has written a *set*, with holes, because completion order does not respect position. If you restore order at the sink, you regain the prefix property; if you do not, resumption has to be driven by which keys already carry a result, not by how far the run got. Decide that before the first production failure rather than during it. ## Saying it well in an interview Name the cause in one sentence — results converge in completion order because durations differ — then ask the question a strong candidate asks: does the sink actually need positional order? Only if the answer is yes do you describe carrying the index, the held-results buffer, the window that bounds it, and what the pipeline does with a straggler.
- What bounds the memory a reordering step needs?Nothing automatic. The fan-out bound caps records being worked on, not finished results waiting behind a straggler: as slots free, new records start and complete, so held results grow with the input. The step needs an explicit window — how far ahead of the oldest unfinished position it may run — plus a policy to stall or fail when that window is exceeded.
- When can you skip reordering entirely?When the sink does not depend on arrival order — each result written against its own key, or inserted into a store that is later read with an explicit sort. Order restoration is needed only where the output position itself carries meaning, such as a positional file or a log consumed by index.
- How does parallelism change how you resume a failed batch?A serial run writes a prefix, so resuming means continuing from the last position. A parallel run writes a set with holes, because completion order ignores position. Either restore order at the sink to regain the prefix property, or drive resumption from which keys already carry a result rather than from how far the run reached.
Numbered pages handed to several typists come back in whatever order each finishes, so you re-collate by the page number written on them rather than by the order they land on the desk.
saying these in an interview costs you the question
- Claims the rejoin preserves source order automatically
- Thinks the concurrency bound also bounds held results
- Says buffering everything and sorting at the end is free
- Believes order can be recovered without carrying a position
- Confuses ordered start with ordered completion
- Ignores what a straggler does to the reordering step