skip to content

How does Flow handle backpressure without an explicit request(n) mechanism like Reactive Streams?

level: seniorimportance: should knowfreq 60%

answer

  1. emit is suspend ⇒ producer waits for collector
  2. Default: same coroutine, sequential, no buffer
  3. Suspension replaces request(n) demand protocol
  4. Opt into concurrency: buffer / conflate / collectLatest
  5. flowOn also adds a channel boundary

basics

~20 s

When the collector is slow, the producer's emit call simply suspends and waits. No values are dropped or buffered by default — the producer politely pauses until the collector is ready for the next value.

solid answer

~40 s

Flow uses suspension as built-in backpressure. emit is a suspend function; by default the producer and collector run sequentially in the same coroutine, so emit doesn't return until the collector has finished processing the previous value. A slow collector therefore naturally throttles a fast producer — the producer suspends at emit rather than overflowing a buffer. This replaces the Reactive Streams request(n) protocol with cooperative suspension. When you want producer/collector to run concurrently, you opt in with buffer() (a bounded channel between them), conflate() (keep only the latest), or collectLatest (cancel slow processing). Those operators introduce capacity and overflow strategies explicitly, so backpressure semantics stay visible rather than hidden in a demand protocol.

code

kotlin · 9 lines
kotlin
// Decouple producer/collector with an explicit bounded buffer:
fastProducer
    .buffer(capacity = 64, onBufferOverflow = BufferOverflow.DROP_OLDEST)
    .collect { slowProcess(it) }

// Or keep only the freshest value when the collector lags:
sensorFlow
    .conflate()
    .collect { render(it) }

go deeper

for a junior

Knows that a slow collector slows the producer and values aren't lost by default.

for a middle

Explains that emit suspends and producer/collector share a coroutine sequentially.

for a senior

Contrasts suspension-based backpressure with request(n), and chooses buffer/conflate/collectLatest with correct overflow semantics.

for a principal

Reasons about capacity, latency, and loss trade-offs across a pipeline, knows flowOn's channel boundary, and designs overflow policy for throughput vs freshness vs memory.

## The problem backpressure solves **Backpressure** is what happens when a producer generates values faster than the consumer can handle them. Naive systems overflow memory or drop data. Reactive Streams (the RxJava/Reactor spec) solves it with an explicit **`request(n)`** demand protocol: the subscriber tells the publisher how many items it can accept. ## How Flow does it instead: suspension In a plain `Flow`, the producer and collector run **sequentially in the same coroutine**. `emit` is a `suspend` function. When the producer calls `emit(x)`, control transfers downstream; `emit` does **not return** until the collector's block for `x` completes. So: - Fast producer + slow collector ⇒ the producer **suspends at `emit`** until the collector is ready. - No unbounded buffer, no dropped values, no `request(n)` accounting. This is *cooperative, demand-free backpressure*: the consumer's pace dictates the producer's pace simply because they share a coroutine and suspension is the hand-off. ```kotlin flow { repeat(3) { i -> println("emit $i") emit(i) // suspends here until collector finishes i } }.collect { i -> delay(100) // slow consumer println("got $i") } // Output interleaves: emit 0, got 0, emit 1, got 1, ... (never races ahead) ``` ## Opting into concurrency / decoupling When you *want* the producer to run ahead of the collector, you make capacity explicit: - **`buffer(capacity, onBufferOverflow)`** — inserts a bounded channel so producer and collector run in **separate coroutines**; producer suspends only when the buffer is full. Overflow strategies: `SUSPEND` (default), `DROP_OLDEST`, `DROP_LATEST`. - **`conflate()`** — buffer of effectively the latest value only; intermediate values are dropped when the collector is slow. - **`collectLatest { }`** — start processing each value but **cancel** the in-flight block if a newer value arrives. - **`flowOn(dispatcher)`** — also introduces a channel boundary (it must, to switch context), giving similar decoupling upstream. ## Why this design - Backpressure becomes a **property of structured concurrency**: suspension already exists for coroutines, so Flow reuses it rather than inventing a demand protocol. - Overflow behavior is **explicit and local** (you see `buffer`/`conflate` in the pipeline) instead of implicit in `request(n)` bookkeeping. - The default (sequential, suspending) is the **safest**: bounded memory, no loss.

  • What's the difference between conflate() and collectLatest in handling a slow collector?
    conflate() drops intermediate values but lets the current collector block finish; collectLatest cancels the in-flight block when a newer value arrives and restarts with the latest.
  • Why does flowOn introduce backpressure decoupling even though you only asked for a dispatcher change?
    Switching context requires handing values across coroutines via a channel, which acts as a buffer/boundary between upstream and downstream.
  • What is the default overflow strategy of buffer()?
    BufferOverflow.SUSPEND — the producer suspends when the buffer is full, preserving backpressure.

A single-lane checkout: the cashier (collector) scans one item before you (producer) can put the next on the belt — the line self-regulates without anyone announcing how many items they'll take.

saying these in an interview costs you the question

  • Saying Flow has no backpressure handling at all
  • Claiming Flow uses request(n) like Reactive Streams internally
  • Believing the producer always runs ahead and buffers values by default
  • Thinking buffer() and conflate() are the same thing
  • Assuming collectLatest waits for the slow block instead of cancelling it

context