skip to content

One stage of a channel pipeline is the bottleneck. How would you fan the work out to parallel workers and fan the results back into a single stream, and what must you decide about ordering, buffer sizes, and shutdown?

level: principalimportance: should knowfreq 33%

answer

  1. N workers on one input channel = free load balancing
  2. fan-in = forwarder per source + barrier + single close
  3. fan-out destroys order: tag+reorder, or partition by key
  4. shallow buffers; utilization near 1 = unstable latency
  5. cancel arm everywhere; close merged output once

basics

~20 s

Fan-out: run N copies of the stage all receiving from the same input channel, so idle workers naturally take the next item. Fan-in: a merger that forwards each worker's output into one channel and closes it after all workers finish. Fan-out loses ordering; keep buffers small and give every worker a cancel path.

solid answer

~1 min

Fan-out is N workers receiving from one shared input channel. Load balancing is automatic — whichever worker is free receives next — so no dispatcher is needed and no worker can be starved by a long item assigned to it. Fan-in merges the N output channels into one: spawn a forwarder per worker, and close the merged channel once all forwarders have finished. That is the same completion barrier used for any multi-producer channel. The three decisions: - **Ordering.** Fan-out destroys input order. If order matters, either tag items with a sequence number and reorder downstream (which costs a buffer bounded by the worst-case skew, and can stall on one slow item), or partition by key so each key goes to one worker and per-key order is preserved. - **Buffers.** Keep them small — roughly a burst, or a small multiple of N. Deep queues raise latency and hide the fact that downstream is the real bottleneck. - **Shutdown.** Every worker selects over its work and a cancel channel; the first error cancels the rest; the merger closes only after all workers are done, otherwise you strand blocked senders. And verify the stage is actually the bottleneck: parallelizing before the constraint just moves the queue.

code

text · 17 lines
text
in     = channel(cap = N)      # shallow
out    = channel(cap = N)
cancel = channel(0)            # closed once to broadcast stop

for i in 1..N:
  spawn worker:
    loop:
      select:
        case item = receive(in):        # closed -> exit loop
          r = process(item)
          select:
            case send(out, r): pass
            case receive(cancel): return
        case receive(cancel): return
    signal(done)

spawn: wait_for(done, times = N); close(out)   # exactly one closer

go deeper

for a junior

Describe the shape: several workers reading one channel, results merged into one channel; mention that results can come back out of order.

for a middle

Add the fan-in mechanics — a forwarder per worker and one close after all finish — and that a bounded input channel preserves backpressure.

for a senior

Lead with ordering strategy, cancellation on every blocking point, and evidence for which stage is actually the bottleneck.

for a principal

Reason quantitatively: concurrency from rate times service time, latency instability as utilization approaches one, shallow buffers as a deliberate choice about where pressure becomes visible, and partition-by-key versus reorder-buffer as an architectural trade.

## The topology Fan-out means several tasks receive from the *same* channel. This is worth pausing on: with one shared input, distribution is a consequence of the channel semantics rather than a scheduling algorithm you write. A worker that is busy is not receiving, so the next item goes to one that is free. That gives work-conserving, self-balancing distribution with no queue-per-worker skew — the failure mode of hash-assigned queues where one worker's queue backs up while others idle. Fan-in is the inverse: N outputs, one stream. It is not free, because a channel's end-of-stream must be closed exactly once by the sending side. The standard construction is one forwarder task per source, all sending into the merged channel, plus a completion barrier: when all forwarders have returned, one coordinator closes the merged channel. in ---> [w1] --\ ---> [w2] ---> merge ---> out ---> [w3] --/ ## Ordering This is the decision people skip. The instant work is distributed across N workers with variable service times, output order stops matching input order. Three coherent answers: 1. **Do not care.** Many pipelines (independent enrichments, writes to an unordered sink) are genuinely order-free. Say so explicitly rather than by accident. 2. **Restore order.** Tag each item with a monotonic sequence number and put a reordering buffer downstream that emits in sequence. This works, but it re-couples the system: the buffer must hold every completed item newer than the oldest outstanding one, so one slow item stalls emission and inflates memory. Bound the buffer and decide what happens when it fills — that limit is your effective skew tolerance. 3. **Partition instead of fan out.** Route by key so all items for a key go to the same worker. Per-key order is preserved without any reorder buffer, and global order is simply not offered. The cost is losing automatic balancing: a hot key can saturate one worker while others idle. This is the same trade partitioned streaming systems make. ## Sizing and buffers How many workers is workload-shaped and should be measured, but the theory frames the target. By Little's law, the average number of items in the system equals arrival rate times average time in system, so to sustain a rate with a given per-item service time you need concurrency at least rate times service time. That number is the *concurrency* needed, which for CPU-bound work is capped by cores and for wait-bound work can far exceed them. Queueing theory supplies the second lesson: as utilization approaches one, waiting time grows without bound. So a pipeline deliberately run at very high utilization has unstable latency, and adding buffer depth does not help — it converts the instability into longer queues. Keep channels shallow (a burst, or a small multiple of N) so that pressure is felt as blocking at the producer rather than hidden as queue depth, and so that a queued item's age stays bounded. Third: parallelizing a stage that is not the constraint moves the queue rather than removing it. Measure occupancy of each channel first — the bottleneck is the stage whose input channel is persistently full and whose output channel is persistently empty. ## Shutdown and failure A fan-out/fan-in section has many blocking points, so it is where leaked tasks accumulate. The invariants: - Every worker's receive and send appears in a select that also watches a cancel channel, so no worker can be stranded by an abandoned consumer. - Cancellation is broadcast by closing the cancel channel once, from the owner of the section. - The merged channel is closed exactly once, by the coordinator, after all forwarders have returned — never by a worker. - Errors travel in-band (an item is either a value or a failure) or on a dedicated channel; the first failure typically triggers cancellation and the section then drains. - If a consumer abandons the merged output, someone must cancel, or the workers block forever on send. A scope that spawns the workers, waits for all of them, and cancels them on exit makes most of this structural rather than a matter of discipline. ## What to say in an interview Lead with the topology, then immediately raise ordering, backpressure and shutdown, because those are where real implementations fail. Finish with measurement: identify the constrained stage, size concurrency from service time and target rate, keep buffers shallow, and re-measure, since the bottleneck moves once you fix one.

  • How do you decide how many workers to fan out to?
    Start from the target rate and measured per-item service time: required concurrency is roughly rate times service time, which for wait-dominated work can be far above the core count and for CPU-dominated work is capped near it. Then measure, because the bottleneck moves — if the downstream channel is persistently full, more workers only deepen queues.
  • The output must stay in input order. What are you signing up for?
    Either a reordering buffer keyed by sequence number, whose size is the tolerated skew and which stalls emission behind one slow item, or partitioning by key so each key is handled by a single worker, which preserves per-key order but gives up automatic load balancing and exposes you to hot keys. Both are real answers; picking neither and hoping order survives fan-out is not.

saying these in an interview costs you the question

  • Assuming output order matches input order after fanning work out to parallel workers.
  • Letting each worker close the merged output channel instead of a single coordinator after a completion barrier.
  • Adding deep buffers between stages to make a pipeline faster — it mainly adds latency and hides the real bottleneck.
  • Fanning out a stage that is not the constraint, which just relocates the queue.
  • Building a per-worker dispatcher when a single shared input channel already balances load.

context