skip to content

How does zip handle concurrency and backpressure between its two source flows, and what are the implications?

level: seniorimportance: should knowfreq 35%

answer

  1. Both sides collected concurrently, joined in lockstep
  2. Fast side runs one step ahead then suspends
  3. Slower source dictates throughput
  4. Either completing cancels the other; extras dropped
  5. Different rates + want-all -> prefer combine

basics

~10 s

zip collects both flows concurrently but emits a pair only when both have produced the next item. The faster flow waits for the slower one, so the slow side sets the pace.

solid answer

~40 s

zip runs both upstreams concurrently (each in its own coroutine) but enforces lockstep at the join point: to emit pair n it needs item n from both sides, so whichever side is faster suspends until the slower side catches up. This couples the pace of the two flows — the slower one dominates throughput. Side effects (emissions, upstream work) on the fast side still run ahead one step but then block. If either side completes, zip stops and cancels the other (structured concurrency); a thrown exception cancels the sibling and propagates. Because of the lockstep coupling, zip is unsuitable when sources have very different rates and you want every value — for that, combine (latest-of-each) or buffering upstream is more appropriate. zip is correct when indices must stay aligned.

code

kotlin · 9 lines
kotlin
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*

val fast = flow { repeat(5) { emit(it) } }
val slow = flow { repeat(5) { delay(100); emit(('a' + it)) } }

// Cadence is set by `slow`; fast suspends after running one step ahead.
fast.zip(slow) { n, c -> "$n$c" }.collect { println(it) }
// 0a, 1b, 2c, 3d, 4e -- roughly one per 100ms

go deeper

for a junior

Knows zip waits for both sides before emitting and that the slow side limits speed.

for a middle

Explains concurrent collection plus lockstep join, and that either flow completing ends the zip.

for a senior

Discusses backpressure coupling, structured-concurrency cancellation/error propagation, and when combine is the better fit.

for a principal

Reasons about throughput characteristics, when zip's coupling is a liability under skewed rates, and the limits of buffer() in changing pairing semantics.

## zip's execution model `zip` is not naive sequential collection. It launches a **separate coroutine to collect the second flow** while the first is collected on the main path, so both upstreams make progress **concurrently**. The synchronization happens at the **join**: emitting pair *n* requires the *n*-th element from **both** sides. ## Backpressure / lockstep coupling Because both sides are needed for each pair: - The **faster** flow can run **one step ahead**, then **suspends** waiting for the slower flow's matching element. - Effective throughput is bounded by the **slower** source — the slow side sets the pace. - This is genuine **backpressure**: zip never buffers an unbounded backlog of the fast side. ```kotlin val fast = flow { repeat(5) { emit(it); /* quick */ } } val slow = flow { repeat(5) { delay(100); emit('a' + it) } } fast.zip(slow) { n, c -> "$n$c" }.collect(::println) // One pair roughly every 100ms — the slow flow dictates the cadence. ``` ## Completion and failure - **Completion:** when **either** flow completes, `zip` completes and **cancels** the other side; unmatched trailing items are dropped. - **Failure:** an exception on one side cancels the sibling coroutine and propagates to the collector — standard **structured concurrency**. - **Cancellation:** cancelling the collector cancels both upstream collections. ## Implications for design - **Use zip** when index alignment is the requirement (pair the *i*-th of two correlated streams) and the rate coupling is acceptable. - **Avoid zip** when sources have very different rates **and** you want every value — the fast side stalls. Reach for **`combine`** (latest-of-each, no stalling on index) or add `buffer()` upstream if you specifically want to decouple producer/consumer for one side (note: buffering changes drop/throughput semantics, not the lockstep pairing itself). - For pure fan-in with no pairing, use **`merge`**. ## Keywords/APIs `Flow.zip`, structured concurrency, `coroutineScope`/child coroutines, backpressure, `buffer`, `combine`, `merge`, cancellation.

  • If the fast flow in a zip is infinite and the slow one emits 3 items, what happens?
    zip emits 3 pairs, then the slow flow completes; zip completes and cancels the infinite flow, so it doesn't run forever.
  • Would adding buffer() to the fast side let zip emit all of the fast flow's extra items?
    No. buffer() can decouple producer/consumer timing, but zip still pairs by index and stops at the shorter flow, so extra items are still dropped.

zip is two runners chained at the wrist: the faster one can stride slightly ahead but must wait for the other before each shared step.

saying these in an interview costs you the question

  • Claiming zip collects the two flows strictly sequentially
  • Saying the fast flow runs unbounded ahead with no backpressure
  • Thinking zip continues after one source completes
  • Believing buffer() removes zip's lockstep pairing
  • Ignoring that one side's failure cancels the other

context