Why does assigning a stream pipeline to a pool of many workers still leave its elements processed one at a time?
answer
- serial by contract, not by accident
- the context picks where, not how many
- one signal at a time, ordered
- operators keep state without locks
- overlap is fanned out explicitly
basics
~20 sA reactive sequence is serial by contract: values reach each stage one at a time, in order, however many workers the execution context owns. That choice decides which worker runs a stage, never how many elements run at once.
solid answer
~50 sAssigning an execution context answers *where* a stage runs, not how much of it runs at once. The sequence itself promises serial delivery: one value is signalled, and the next is signalled only after the first has been handed on. That promise is why an operator can keep unsynchronised state such as a running total or a set of keys already seen. So a scoring stage on a sixteen-worker context still has exactly one record inside it at any instant — one worker busy, fifteen parked — and wall clock stays close to `record count × per-element cost`. Parallelism has to be built explicitly on top of the sequence: fan each element out into its own inner operation and keep several of those running under a concurrency bound, or split the sequence into a fixed number of independent tracks and rejoin them afterwards.
code
pseudocode · 11 lines// one sequence, sixteen workers available
pipeline = source(records)
.map(record -> score(record)) // the whole scoring stage
.executeOn(poolOfSixteenWorkers)
// what happens over time:
// some worker: score(r1)
// some worker: score(r2) -- possibly a different one, never at the same time
// some worker: score(r3)
// exactly one record is inside map at any instant
// wall clock stays about count(records) * cost(score)go deeper
Remember the one-line fact: a stream pipeline processes elements one at a time, and picking a bigger pool of workers does not change that. It chooses where the work runs, not how much runs together.
Explain the mechanism: the signal contract serialises delivery so operators can hold state without locks, and an execution context selects a worker rather than lifting that serialisation. Be able to name what would break if it did not hold.
Show how you would prove it on a real batch: wall clock tracking record count times per-element cost, one hot worker, and a profile that never shows the stage on two stacks at once. Then say which explicit construction you would add.
Frame it as a cost decision. Parallelism inside a pipeline is a construction with permanent costs — a bound to tune, order to restore, harder failure reproduction — so it should be adopted where measurement says the pipeline is the bottleneck, not as a default posture.
## Why a sequence is serial in the first place A source and its subscriber agree on a **signal contract**: values are handed over one at a time, and a value is signalled only after the previous one has been passed on. A terminal signal — completion or failure — arrives last and exactly once. That serialisation is not an implementation quirk of one reactive library; it is what makes the paradigm usable at all. Serial delivery is what lets operators hold ordinary, **unsynchronised** state: - a running total or count accumulated across elements; - a set of keys already seen, for a de-duplicating stage; - a window of recent values waiting to be emitted as a group; - a flag recording that a boundary value has passed. If two values could enter one stage simultaneously, every operator anyone has ever written would need a lock, and the ordering that downstream stages depend on would be gone. So *the sequence is serial* is a promise the pipeline makes to its own stages, and it holds no matter what the work runs on. ## What choosing an execution context actually changes Assigning the pipeline to a context — a pool of many workers, an elastic pool, a single dedicated worker — answers which worker carries the signals and therefore which thread the stage's code executes on. It also decides whether the thread that started the subscription is released instead of being used to run the work. It does not multiply the number of values in flight through a stage. | Assigning a many-worker context decides | It does not decide | |---|---| | Which worker executes a stage's code | How many elements are inside that stage at once | | Whether the subscribing thread is released | Whether per-element work overlaps in time | | That the work is off the caller's thread | How long the run takes per element | One subtlety is worth saying out loud, because interviewers probe it: **serial does not mean the same worker forever**. Across the life of a sequence, successive deliveries may be carried by different workers of the same context, as long as the hand-over is ordered so that state written during one delivery is visible to the next. What the contract forbids is two deliveries *overlapping*. ## The symptom in a nightly batch Take a nightly scoring batch: read a large record set, score each record, write the results back. You point the pipeline at a sixteen-worker context and the run takes exactly as long as before. What you observe: 1. Wall clock stays close to `record count × per-element cost` — the number a plain loop would give. 2. One worker runs hot; the rest are parked with nothing to take. 3. Enlarging the context changes nothing, because what the workers are waiting on never holds more than one piece of real work. 4. A sampling profile never shows the scoring function on two call stacks at the same instant. That fourth observation is the one that settles the argument. Concurrency you cannot see in two simultaneous stacks is concurrency that does not exist. ## Where the parallelism has to come from Real overlap is created, not configured. Two routes exist, and both are explicit constructions on top of the serial sequence: - **Fan out per element.** Turn each record into its own inner operation, keep several of those subscribed at once under a **concurrency bound**, and merge their results back into one sequence. The unit handed out is one element. - **Split into fixed tracks.** Route elements round-robin into a fixed number of tracks, each of which is itself a serial sequence pinned to a worker, then rejoin the tracks into one sequence. The unit handed out is a share of the stream. Both routes take on costs the serial pipeline never had: a bound to choose, and output order that the rejoin no longer preserves. Neither is bought by enlarging a pool. ## When the time is spent waiting, not computing If the per-element cost is a remote call rather than processor work, more workers still do not help. A serial sequence issues one call, waits for it, then issues the next, so exactly one call is outstanding at any moment. Overlap requires several calls in flight, which is the fan-out route again. Extra workers only provide more places to be idle. ## How to answer it in an interview 1. State the contract: one signal at a time, ordered, terminal signal last. 2. Separate *where* from *how many* — the context answers only the first. 3. Name the two honest routes to parallelism and the unit each hands out. 4. Name what you take on by using them: a bound, and restoring order at the rejoin.
- Does serial delivery mean every element is carried by the same worker for the life of the sequence?No. Over time a stage's signals may be carried by different workers of the same context; an elastic or work-stealing context is free to drain the next signals elsewhere. What the contract forbids is two deliveries overlapping, and it requires the hand-over to be ordered so state written during one delivery is visible to the next.
- If the slow stage spends its time waiting on a remote call rather than on the processor, do more workers help?Not by themselves. A serial sequence issues one call, waits, then issues the next, so only one call is ever outstanding and waiting time never overlaps. Overlap comes from having many calls in flight at once, which means fanning out under a bound. Extra workers simply give the pipeline more places to be idle.
saying these in an interview costs you the question
- Says a multi-worker pool automatically parallelises the pipeline
- Thinks switching execution context speeds up per-element work
- Believes more workers means more elements in flight on one sequence
- Claims a stage may receive two values concurrently
- Confuses choosing a worker with creating concurrency