skip to content

What does the buffer() operator do to a Kotlin Flow, and why would you add it to a pipeline?

level: juniorimportance: must knowfreq 70%

answer

  1. Sequential by default: emit suspends until collect finishes
  2. buffer() = separate coroutine + Channel between stages
  3. Overlaps producer and collector, total time -> slower stage
  4. Default SUSPEND loses nothing; order preserved
  5. capacity default Channel.BUFFERED (64)

basics

~20 s

buffer() lets the part producing values and the part collecting them run at the same time instead of taking turns. Producers can race ahead into a small queue, so a slow collector no longer slows the producer.

solid answer

~40 s

By default a Flow runs sequentially: emit() suspends until collect() finishes processing that item, so producer and collector share one coroutine and alternate. buffer() inserts a channel between them and runs the upstream in a separate coroutine, so emission and collection happen concurrently. Emitted values are placed in the buffer; the collector drains them at its own pace. This can shorten total time when both stages take real work, because the producer keeps emitting while the collector processes earlier items. You pass an optional capacity (default Channel.BUFFERED, 64) and an onBufferOverflow policy. buffer() does not lose items by default (BufferOverflow.SUSPEND), unlike conflate(). It changes timing/concurrency, never the values or their order.

code

kotlin · 13 lines
kotlin
val flow = flow {
    repeat(3) { i ->
        delay(100)
        emit(i)
    }
}

// Without buffer: ~1200ms (stages alternate)
// With buffer:    ~1000ms (stages overlap)
flow.buffer().collect { value ->
    delay(300)
    println(value)
}

go deeper

for a junior

Knows buffer() lets producer and collector run concurrently and speeds up pipelines with two slow stages.

for a middle

Explains the default-sequential model (emit suspends on collect) and that buffer inserts a channel + extra coroutine, preserving order and not dropping by default.

for a senior

Discusses capacity defaults, the SUSPEND back-pressure guarantee, and operator fusion with flowOn.

for a principal

Frames buffer() as one point on the back-pressure spectrum and reasons about when overlap actually helps vs. adds latency/memory.

## The problem buffer() solves A cold `Flow` is **sequential by default**. Inside `collect { }`, each `emit(value)` call in the producer **suspends** until the collector's lambda finishes processing that value. Producer and collector therefore run on the **same coroutine** and effectively take turns: emit, process, emit, process. If the producer takes 100 ms per item and the collector takes 300 ms, each item costs 400 ms. ## What buffer() changes `buffer()` decouples the producer from the collector by inserting a **Channel** between them and running the **upstream in its own coroutine**. Now the producer can emit into the buffer while the collector independently drains it. The two stages overlap, so total time approaches the time of the **slower** stage rather than the **sum** of both. ```kotlin flow { repeat(3) { i -> delay(100) // produce emit(i) } }.buffer() // upstream now runs concurrently .collect { value -> delay(300) // consume println(value) } ``` Without `buffer()` this takes ~3 * 400 ms. With `buffer()` the producer races ahead, so it's ~100 + 3 * 300 ms. ## Signature and parameters ```kotlin fun <T> Flow<T>.buffer( capacity: Int = Channel.BUFFERED, // default 64 onBufferOverflow: BufferOverflow = BufferOverflow.SUSPEND ): Flow<T> ``` - **capacity** — how many items the buffer holds. Special values: `Channel.RENDEZVOUS` (0), `Channel.CONFLATED`, `Channel.UNLIMITED`, `Channel.BUFFERED` (default, 64 unless overridden by a system property). - **onBufferOverflow** — what happens when a full buffer meets a new emission: `SUSPEND` (default — back-pressure the producer), `DROP_OLDEST`, or `DROP_LATEST`. ## Key guarantees - **Order preserved**, **values unchanged** — `buffer()` only affects *timing/concurrency*. - With the default `SUSPEND` policy **nothing is lost**: a full buffer suspends the producer until space frees up (true back-pressure). - Adjacent `buffer()` calls **fuse**: `flowOn` and `buffer` combine instead of stacking channels. ## How it differs from siblings - `conflate()` is shorthand for `buffer(CONFLATED)` / effectively `DROP_OLDEST` capacity-1 — it **drops** intermediate values. - `collectLatest`/`mapLatest` **cancel** the in-flight collector when a new value arrives. `buffer()` is the *lossless concurrency* tool; reach for it when both stages do real work and you want them to overlap without dropping data.

  • Does buffer() ever drop or reorder values?
    With the default SUSPEND policy it never drops, and order is always preserved. Only DROP_OLDEST/DROP_LATEST drop, and even then order of surviving items is kept.
  • What is the default capacity?
    Channel.BUFFERED, which is 64 unless overridden by the kotlinx.coroutines.channels.defaultBuffer system property.

Like a conveyor belt with a holding tray between two workers: the first worker stacks parts on the tray and keeps going instead of waiting for the second worker to finish each one.

saying these in an interview costs you the question

  • Claiming buffer() changes or transforms the emitted values
  • Thinking buffer() drops items by default (it suspends)
  • Saying Flow is concurrent by default and buffer is redundant
  • Confusing buffer() with conflate() (lossless vs lossy)

context