skip to content

Within a single collector, how do emit and the collector body interact? Are they concurrent or sequential?

level: middleimportance: should knowfreq 55%

answer

  1. emit suspends until collector finishes
  2. Sequential handshake, no concurrency by default
  3. Automatic backpressure
  4. buffer()/conflate() decouple producer and collector
  5. Concurrent emit needs channelFlow

basics

~10 s

They run one after another, not at the same time. emit suspends the producer until the collector finishes handling the current value, then the producer continues to the next emit.

solid answer

~40 s

In a plain cold Flow, the producer and the single collector run sequentially in the same coroutine flow. Calling emit suspends the producer until the downstream collector lambda finishes processing that value; only then does emit return and the producer proceed to the next emission. This gives natural, built-in backpressure: a slow collector slows the producer. There is no concurrency or buffering by default. If you want the producer to run ahead while the collector works, you insert buffer(), which decouples them via a channel, or conflate() to keep only the latest. Operators like flatMapMerge or flowOn introduce concurrency/context shifts. But out of the box, emit -> collector -> next emit is a strict, ordered, single-threaded-per-context handshake.

code

kotlin · 10 lines
kotlin
flow {
    emit(1); println("producer continues")
    emit(2)
}.collect {
    println("collect $it"); delay(100)
}
// collect 1 -> (waits 100ms) -> producer continues -> collect 2

// Decouple so producer runs ahead:
flow { emit(1); emit(2) }.buffer().collect { delay(100) }

go deeper

for a junior

Knows values arrive in order and one at a time.

for a middle

Explains that emit suspends until the collector finishes, giving backpressure.

for a senior

Knows buffer/conflate/collectLatest/flowOn and how each alters the sequential model and concurrency.

for a principal

Reasons about backpressure strategy, throughput vs latency, and where to introduce buffering in a pipeline.

## Sequential handshake For a single collection of a cold Flow, the producer (the `flow {}` block) and the collector lambda execute **sequentially**, taking turns: 1. Producer calls `emit(x)`. 2. `emit` **suspends** the producer and runs the downstream collector body with `x`. 3. When the collector body returns, `emit` resumes. 4. Producer continues to the next statement / next `emit`. There is **no parallelism** between producing and consuming in a plain flow. This is why emissions arrive in order and why a slow collector throttles the producer — that is **backpressure**, and it is automatic. ```kotlin flow { emit(1) println("after emit 1") emit(2) }.collect { println("got $it") delay(100) } // got 1 -> (100ms) -> after emit 1 -> got 2 ``` Note `after emit 1` prints **after** the collector finished value 1, proving the suspension. ## Adding concurrency / buffering If you don't want the producer to wait for the collector, you opt in: - `buffer(capacity)` — runs producer and collector in **separate coroutines** connected by a channel, so the producer can emit ahead while the collector processes. Decouples them. - `conflate()` — like buffer but drops intermediate values, keeping only the latest. - `collectLatest { }` — cancels the in-flight collector body when a new value arrives. - `flowOn(dispatcher)` — changes the **upstream** dispatcher; the producer can then run on a different thread than the collector, but emit/collect ordering within the stream is still preserved. - `flatMapMerge` / `channelFlow` — introduce real concurrency for multiple inner flows. ## Why sequential by default Keeping emit/collect sequential makes flows **predictable and ordered**, gives free backpressure, and avoids hidden threads. Concurrency is always an explicit operator, never an accident. ## Common gotcha Launching a new coroutine inside `flow {}` to call `emit` from it throws `IllegalStateException` (emit is not thread-safe across coroutines). Use `channelFlow {}` + `send` when you genuinely need concurrent emission.

  • What does buffer() change about emit/collect timing?
    It moves the collector to a separate coroutine via a channel, so the producer can keep emitting without waiting for the collector — decoupling them.
  • Why does calling emit from a launched child coroutine inside flow {} fail?
    Flow's emit must be called from the same coroutine that runs the block to preserve context and sequencing; cross-coroutine emit throws IllegalStateException. Use channelFlow with send instead.

Like passing a single baton: the producer can't grab it again until the collector hands it back.

saying these in an interview costs you the question

  • Saying the producer runs in parallel with the collector by default
  • Thinking flows buffer values automatically
  • Claiming emit returns immediately regardless of the collector
  • Not recognizing emit/collect as the source of backpressure

context