skip to content

How does buffer() interact with flowOn, and what does 'operator fusion' mean for a chain like flowOn(...).buffer().map { }?

level: seniorimportance: should knowfreq 30%

answer

  1. buffer & flowOn both add a channel + upstream coroutine
  2. Adjacent buffer/flowOn/conflate FUSE into one channel
  3. Later buffer overrides flowOn's implicit capacity
  4. Transform between them splits the boundary
  5. Placement decides which coroutine runs each map

basics

~20 s

flowOn already adds a buffer to move work to another thread, and the coroutines library merges adjacent buffer/flowOn stages into one channel instead of stacking several. So extra buffer() calls next to flowOn usually just tune that single buffer rather than adding new ones.

solid answer

~40 s

Both `buffer()` and `flowOn()` introduce a channel to run upstream on a separate coroutine (flowOn additionally switches the `CoroutineContext`/dispatcher). The Flow machinery performs **operator fusion**: directly adjacent `buffer`/`flowOn`/`conflate`/`channelFlow` boundaries are combined into a **single** channel with a single coroutine, so you don't pay for multiple hops. When you write `upstream.flowOn(IO).buffer(N)`, the buffer request and flowOn boundary fuse; the effective capacity/overflow is reconciled (a later explicit `buffer` can override the implicit one flowOn would use). Crucially, position matters relative to operators in between: `map { }` placed before vs after a `buffer()` runs upstream-of vs downstream-of the channel boundary, changing which coroutine executes the transform. Fusion only merges *adjacent* buffering operators; an intervening transform splits them. Ordering and values are still preserved.

go deeper

for a junior

Knows flowOn switches dispatcher and buffer adds concurrency, even if fusion details are fuzzy.

for a middle

Understands both add a channel and that placement of map relative to buffer matters.

for a senior

Explains operator fusion: adjacent buffering operators merge into one channel, and a later buffer overrides flowOn's implicit capacity.

for a principal

Designs pipeline partitioning intentionally, reasoning about coroutine boundaries, allocation cost, and where each transform executes.

## Two operators, one underlying mechanism - `flowOn(context)` changes the `CoroutineContext` (e.g. dispatcher) of **upstream** operators and, to cross the context boundary safely, runs that upstream in a **separate coroutine connected by a channel**. - `buffer(capacity, onBufferOverflow)` runs upstream in a separate coroutine connected by a channel **without** changing context. Both therefore create a producer/consumer boundary backed by a `Channel`. ## Operator fusion To avoid paying for several channels in a row, kotlinx.coroutines **fuses** directly adjacent *fusible* operators — `buffer`, `flowOn`, `conflate`, and the implicit buffering of `channelFlow`/`produce`. Instead of `channel -> channel -> channel`, the runtime collapses them into **one** channel with **one** coroutine and a reconciled capacity/overflow. ```kotlin source .flowOn(Dispatchers.IO) // implies a buffer to cross contexts .buffer(64) // FUSES with the flowOn boundary ``` Here you don't get two channels: the explicit `buffer(64)` *configures* the single fused boundary. A later `buffer` generally overrides the capacity/overflow that flowOn would otherwise pick. ## Position changes which coroutine runs a transform Fusion only merges **adjacent** buffering operators. Put a transform between them and the boundary is real on both sides: ```kotlin source .map { heavyA(it) } // runs UPSTREAM of the boundary (producer coroutine) .buffer() .map { heavyB(it) } // runs DOWNSTREAM (collector coroutine) .collect { ... } ``` `heavyA` and `heavyB` now execute on **different coroutines**, overlapping in time. Moving `buffer()` up or down the chain re-partitions the work. This is the main reason buffer placement is a deliberate design choice, not noise. ## What fusion does NOT change - **Order** of values and the **values themselves** are untouched. - It's purely a performance/structure optimization: fewer channel allocations, fewer coroutine hops. ## Practical guidance - Don't sprinkle redundant `buffer()` right next to `flowOn` expecting additive buffering — they fuse; instead pass the capacity you want. - Use buffer **placement** to decide which stage runs concurrently with which. - Verify with reasoning about coroutines, not by assuming each operator adds latency. ```kotlin // Effective: one fused boundary, IO upstream, 128-slot back-pressured buffer val result = source .map { parse(it) } .flowOn(Dispatchers.IO) .buffer(128) .map { render(it) } ```

  • Does putting buffer() right after flowOn() create two channels?
    No. They fuse into a single channel/coroutine; the explicit buffer just configures that one boundary's capacity and overflow policy.
  • If you place map { } before vs after buffer(), what differs?
    The map before runs on the producer coroutine (upstream of the channel); the map after runs on the collector coroutine, so the two transforms can overlap.

saying these in an interview costs you the question

  • Believing buffer next to flowOn stacks multiple channels
  • Thinking fusion can reorder or drop values
  • Assuming buffer placement is irrelevant to where work runs
  • Claiming flowOn does not involve any buffering

context