A nightly scoring batch must use every core, so what are the two ways a serial reactive sequence gets real parallelism?
answer
- two constructions, not a setting
- unit of work: element versus share
- inner operation must leave the driving thread
- tracks assigned up front, cannot rebalance
- both lose source order at the rejoin
basics
~20 sTwo routes exist. Turn each element into an inner operation carrying its own work and keep several subscribed at once under a concurrency bound; or split the sequence into a fixed number of worker-pinned tracks and rejoin them afterwards.
solid answer
~50 sThe first route fans out **per element**: each record becomes its own inner operation, and the pipeline keeps up to a bounded number of them running at once, merging results as they finish. The second route splits the sequence into a **fixed number of tracks** — typically round-robin by position — where each track is its own serial sequence pinned to a worker, and the tracks are rejoined into one sequence at the end. The trap in the first route is that the inner operation must actually be moved onto a worker or be genuinely asynchronous; a synchronous inner operation runs during subscription, on the very thread driving the outer sequence, so nothing overlaps. Fan-out absorbs uneven per-record costs because it refills from a shared pool of pending work; fixed tracks pay fewer hand-offs but cannot rebalance a track that drew the expensive records.
code
pseudocode · 11 linesmaxInFlight = 8
results = source(records)
.fanOut(
record -> innerOperation(() -> score(record))
.executeOn(workerPool), // <- without this, no overlap
maxInFlight) // at most 8 inner operations at once
.subscribe(write)
// drop .executeOn and the inner operation runs during subscription,
// on the thread driving source(records) -- one record at a timego deeper
Learn the two names and the unit each hands out: one element at a time under a limit, or a fixed set of tracks each taking a share of the stream. Both are things you build, not settings you switch on.
Explain the mechanics of each and the trap in the first: a synchronous inner operation runs during subscription on the driving thread, so the fan-out overlaps nothing. Be ready to say which shape suits waiting work and which suits uniform processor work.
Argue the choice from measured per-element cost and its variance, then account for what both routes cost you afterwards: a bound, lost source order, and a dependency that now sees several concurrent callers.
Own the default. Decide which construction your batch pipelines standardise on, what evidence justifies deviating, and how the bound is chosen and reviewed, so that every team does not rediscover skew and unbounded fan-out separately.
## The starting point A sequence is serial by contract, so both routes are constructions layered on top of it. Both take the same two costs: a **bound** to choose, and **source order** that the rejoin no longer preserves. What differs is the unit of work handed to a worker. ## Route one — fan out per element Each element is mapped to an **inner operation** that performs the work for that element and produces its result. The pipeline subscribes to several inner operations at once, up to a **concurrency bound**, and merges their results into the outer sequence as each completes. When one finishes, the next pending element starts. What this buys, and what it demands: - The unit is one element, so uneven per-record costs even out: a slow record occupies one slot while the others keep cycling. - The bound is explicit, so peak in-flight work does not grow with the size of the input. - It fits work that *waits* — remote calls, queued requests — because many can be outstanding at once. - **The inner operation must genuinely leave the driving thread.** If it is synchronous, subscribing to it performs the work then and there, on the thread advancing the outer sequence. The fan-out hands out work and then does it itself, one element at a time, and the measurement looks exactly like no change at all. That last point is the single most common way a fan-out is written and produces nothing. It is worth stating in an interview before anyone asks. ## Route two — split into fixed tracks The sequence is partitioned into a fixed number of **tracks**. Elements are distributed across them by position, usually round-robin; each track is itself a serial sequence assigned to a worker; the stage runs inside each track; the tracks are then rejoined into one sequence. - The unit is a share of the whole stream, not a single element. - Ordering *within* one track is preserved; ordering across the rejoin is not. - Hand-offs are paid per element per track rather than by constructing and subscribing to an inner operation for every record. - Assignment happens up front, so it cannot rebalance: a track that drew the expensive records stays busy while the others idle and the rejoin waits on the slowest. ## Choosing between them | | Fan out per element | Fixed tracks | |---|---|---| | Unit handed to a worker | One element | A share of the stream | | Fits | Waiting work, uneven per-element cost | Uniform processor-bound work | | Rebalancing | Refills from pending work as slots free | Fixed at assignment, cannot rebalance | | Per-element overhead | Constructing and subscribing an inner operation | One routing hop into a track | | Natural bound | Explicit maximum in flight | The track count | A practical rule: if the per-element cost is dominated by waiting, fan out; if it is uniform processor work whose per-element duration is small relative to a hand-off, tracks win. If the costs are wildly uneven, fan out regardless, because skew is what tracks handle worst. ## What neither route gives you Neither route makes the *outer* sequence non-serial. Downstream of the merge or the rejoin you are back to one value at a time, which is what lets you write an ordinary aggregating or writing stage after it. Both routes also lose source order at the point where results converge, so a batch whose output position carries meaning must restore it deliberately. Neither route fixes a pipeline whose time is spent somewhere else either. If the write-back is the bottleneck, scoring eight records at a time simply makes eight records wait at the write. ## A worked sizing sketch Suppose a million records, each scored in about ten milliseconds of processor work, on a machine with eight usable cores. 1. Serially, the scoring alone is about `1,000,000 × 10 ms`, roughly three hours. 2. With eight-way parallelism the floor is about a quarter of an hour, minus whatever the rest of the pipeline costs. 3. Because the work is uniform and processor-bound, fixed tracks are the cheaper construction: eight tracks, one hand-off per record, no inner operation per record. 4. If scoring instead meant a remote call of about ten milliseconds, tracks would leave every worker blocked on a socket, and a fan-out with a bound chosen from what the dependency can absorb would be the right shape. Stating that reasoning out loud — unit of work, shape of the cost, which construction follows — is what separates an answer that names both routes from one that has used them.
- Why does mapping each element to an inner operation sometimes give no speedup at all?Because the inner operation is synchronous. Subscribing to it performs the work immediately, on the thread driving the outer sequence, so the fan-out hands out work it then executes itself, one element at a time. The inner operation has to be placed on a worker, or be genuinely asynchronous, before anything overlaps.
- When would you choose fixed tracks over per-element fan-out?When the work is uniform and processor-bound and the per-element cost is small next to the cost of constructing and subscribing an inner operation for every record. Tracks pay one routing hop instead. Uneven per-element cost argues the other way, because a track that drew the expensive records cannot be rebalanced and the rejoin waits on it.
- Does either route make the sequence after the merge or rejoin non-serial?No. Downstream of the convergence point you are back to one value at a time, which is what lets an ordinary aggregating or writing stage follow it safely. The parallelism lives strictly between the fan-out and the merge, or inside the tracks before the rejoin.
saying these in an interview costs you the question
- Thinks a bigger worker pool is one of the two routes
- Fans out without moving the inner operation off the driving thread
- Believes rejoining tracks restores the original order
- Says every element should get its own track
- Treats an unbounded fan-out as the default shape
- Picks fixed tracks for work that mostly waits on a dependency