skip to content

buffer & Overflow Strategies

buffer puts producer and collector in separate coroutines with a queue between them, and the overflow policy decides whether to suspend or drop. It is the first tool to reach for when a slow collector is throttling a fast upstream.

part ofKotlinoverview, primer and where to startread it →
on this pageshow

questions

5

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

open as a page

Explain the three BufferOverflow strategies for buffer() and what each does when the buffer is full.

level: middleimportance: must knowfreq 60%

basics

~20 s

When the buffer fills up: SUSPEND pauses the producer until there is room (nothing lost); DROP_OLDEST throws away the oldest queued item to make space for the new one; DROP_LATEST throws away the new item and keeps what is already queued.

open as a page

What capacity does buffer() use by default, and what do the special Channel capacity constants mean in this context?

level: middleimportance: should knowfreq 40%

basics

~20 s

By default buffer() holds about 64 items. You can also pass special values: 0 for no buffer (hand-off), a large unlimited buffer that never blocks, or a 1-slot conflated buffer that always keeps the newest value.

open as a page

You have a fast event producer and a collector that can fall behind (e.g. rendering live telemetry). How do you reason about buffer capacity and BufferOverflow choice, and what are the failure modes of each?

level: seniorimportance: should knowfreq 28%

basics

~20 s

Decide whether you must keep every event or only the latest. If every event matters, suspend or use a big bounded buffer. If only freshness matters, drop oldest with a tiny buffer. Avoid unlimited buffers because they can run out of memory.

open as a page

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

level: seniorimportance: should knowfreq 30%

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.

open as a page