skip to content

Backpressure

What to do when the producer outruns the collector: buffer, conflate, drop to the latest value, or let suspension slow the producer down. Interviewers reach for this whenever the scenario involves sensor data or a fast event source.

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

explore

questions

15

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

What does Flow's collectLatest do, and how is it different from a plain collect?

level: juniorimportance: must knowfreq 70%

basics

~20 s

collectLatest runs your handling code for each value, but if a new value arrives before the previous one finishes, it stops the previous work and starts on the new value. Plain collect always finishes every value first.

open as a page

What does the conflate() operator do to a Kotlin Flow, and when would you reach for it?

level: juniorimportance: must knowfreq 55%

basics

~20 s

conflate() lets a fast producer skip ahead: if the collector is busy, in-between values are dropped and it only gets the newest one. Good when you only care about the latest value, like a live counter.

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

Why might collectLatest fail to cancel an in-flight block, and how do you make CPU-bound work in the block actually cancellable?

level: middleimportance: must knowfreq 55%

basics

~20 s

Cancellation only happens when the running code reaches a pause point. If your block does heavy non-stop computation, there's no pause point, so the old block keeps running. Add checks like ensureActive() or yield() so it can stop.

open as a page

conflate() is equivalent to which buffer(...) configuration, and what does that imply about its buffer size and overflow handling?

level: middleimportance: must knowfreq 45%

basics

~10 s

conflate() is the same as a buffer of size one that throws away the oldest value when a new one arrives. So at most one value waits, and it's always the newest.

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

Compare mapLatest and transformLatest with map/transform. What exactly gets cancelled, and how do they relate to collectLatest?

level: middleimportance: should knowfreq 45%

basics

~20 s

map and transform run their lambda fully for each value. mapLatest and transformLatest start the lambda for each value but cancel the previous lambda when a new value comes in. They are the same 'latest wins' idea as collectLatest, but as intermediate operators that produce a new Flow.

open as a page

Given a producer emitting 1..100 every 1ms through conflate() into a collector that takes 100ms per item, roughly how many items get collected and which ones?

level: middleimportance: should knowfreq 35%

basics

~10 s

Far fewer than 100 — roughly one item per 100ms of collector time. You'll see a sparse, increasing set ending at the last value (100), with most in-between numbers dropped.

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

You are building a live-search feature. Compare collectLatest/flatMapLatest with conflate and explain how debounce fits in. When is collectLatest the wrong choice?

level: seniorimportance: should knowfreq 40%

basics

~20 s

For live search you want only the newest query to matter. collectLatest/flatMapLatest cancel the previous in-flight search when a new query arrives. conflate instead keeps the latest value but lets the current handler finish without cancelling. debounce waits for a pause in typing before emitting, cutting wasted requests.

open as a page

Both conflate() and collectLatest help a slow consumer keep up with a fast producer. How do they differ in what they do with the in-progress work?

level: seniorimportance: should knowfreq 40%

basics

~10 s

conflate() lets the current item finish, then jumps to the newest waiting value. collectLatest cancels the current work as soon as a newer value arrives and restarts with it. One finishes-then-skips, the other interrupts.

open as a page

Where in a Flow chain does conflate() take effect, and what determines which dispatcher/coroutine the dropped-or-kept boundary sits on?

level: seniorimportance: nice to knowfreq 22%

basics

~20 s

conflate() splits the chain at the point you call it: everything before it runs in one coroutine, everything after in another, and dropping happens at that split. Operators upstream of conflate keep running; only values crossing the split get conflated.

open as a page

Explain the structured-concurrency mechanism behind collectLatest's cancel-and-restart, and the correctness and performance tradeoffs of relying on it at scale.

level: principalimportance: nice to knowfreq 20%

basics

~20 s

collectLatest runs each value's handler in a child coroutine of the collector. When a new value arrives it cancels that child and launches a fresh one. So work is cheap to throw away, but cancelling has costs and only completed work has real effects — you must design handlers to be safe if abandoned partway.

open as a page