skip to content

How does buffering and backpressure work in channelFlow, and what is the default channel capacity?

level: seniorimportance: should knowfreq 38%

answer

  1. Default capacity = RENDEZVOUS (0)
  2. send() suspends => natural backpressure
  3. buffer(capacity, onBufferOverflow) tunes it
  4. CONFLATED / conflate() = latest only
  5. buffer/conflate/flowOn are fused

basics

~10 s

By default the backing channel has no buffer (rendezvous), so send() suspends until the collector takes the value, giving natural backpressure. Adding buffer() or a capacity lets producers run ahead.

solid answer

~40 s

channelFlow's backing channel defaults to RENDEZVOUS capacity (0): each send() suspends until the collector is ready to receive, so fast producers are throttled by a slow collector — backpressure for free. You change this by applying the buffer(capacity) operator downstream, which sets the channel's buffer size (e.g., buffer(64), or Channel.UNLIMITED, Channel.CONFLATED). With a buffer, producers can send up to capacity values before suspending. onBufferOverflow controls behavior when full: SUSPEND (default), DROP_OLDEST, or DROP_LATEST. conflate() is shorthand for keeping only the latest. Because send() is a suspend function, backpressure naturally propagates back to every producing coroutine. trySend() is the non-suspending variant that returns a ChannelResult and is useful in callback contexts where you can't suspend. Fusion: consecutive buffer()/flowOn operators are fused so you don't stack redundant channels.

code

kotlin · 15 lines
kotlin
import kotlinx.coroutines.channels.*
import kotlinx.coroutines.flow.*

val events = channelFlow {
    repeat(10_000) { send(it) }
}
// keep only newest, never block producer
.buffer(capacity = 1, onBufferOverflow = BufferOverflow.DROP_OLDEST)

// non-suspending emit from a non-suspend context:
val cb = channelFlow<Int> {
    val r: ChannelResult<Unit> = trySend(42)
    if (r.isFailure) { /* dropped or full */ }
    awaitClose { }
}

go deeper

for a junior

Knows send() can suspend and that there is some notion of backpressure.

for a middle

States the rendezvous default and that buffer() changes capacity and overflow behavior.

for a senior

Explains overflow strategies, conflate, trySend vs send, and operator fusion with concrete trade-offs.

for a principal

Designs throughput-vs-memory policy across a pipeline, choosing capacities/overflow per stage and reasoning about UNLIMITED OOM risk.

## The backing channel and its default capacity Every `channelFlow` is backed by a `Channel`. Its **default capacity is `RENDEZVOUS` (0)**: there is no buffer, so a `send(value)` call **suspends until the collector is ready to receive** that value. This gives **backpressure** automatically — a slow collector slows down the producers. ## Tuning the buffer You rarely set capacity on the channel directly; instead you apply the **`buffer()`** operator downstream, which configures the upstream channel: ```kotlin channelFlow { repeat(1000) { send(it) } } .buffer(capacity = 64) // producer may run up to 64 ahead .collect { slowConsume(it) } ``` Useful capacities: - **`Channel.RENDEZVOUS` (0)** — default, strict handoff. - **`Channel.BUFFERED`** — default buffered size (64 unless overridden). - **`Channel.UNLIMITED`** — never suspends on send (risk of OOM). - **`Channel.CONFLATED`** — keep only the latest value (equivalent to `conflate()`). ## Overflow strategy `buffer(capacity, onBufferOverflow = ...)` accepts a `BufferOverflow`: - **`SUSPEND`** (default) — producer suspends when full. - **`DROP_OLDEST`** — discard the oldest buffered value. - **`DROP_LATEST`** — discard the incoming value. With `DROP_OLDEST`/`DROP_LATEST`, `send()` no longer suspends on overflow — useful when you must never block producers (e.g., UI events). ## suspend vs. trySend - **`send()`** — `suspend`; respects backpressure. - **`trySend()`** — non-suspending; returns a `ChannelResult<Unit>` (`isSuccess`/`isFailure`/`isClosed`). Handy inside non-suspending callbacks. ## Operator fusion The flow machinery **fuses** adjacent `buffer`, `conflate`, and `flowOn` operators so you don't pay for multiple chained channels — the effective buffer/dispatcher is computed once. This means stacking `.buffer().buffer()` doesn't create two channels; the most relevant configuration wins/combines per documented rules. ## Practical guidance - Keep the rendezvous default when correctness depends on not overrunning the consumer. - Add `buffer(n)` to overlap producer and consumer work for throughput. - Use `conflate()`/`CONFLATED` for latest-wins UI state where stale intermediates can be dropped.

  • If you set Channel.UNLIMITED, what risk do you introduce?
    send() never suspends, so a fast producer with a slow collector can accumulate unbounded values in memory and cause OutOfMemoryError. You trade backpressure safety for throughput.
  • What is the difference between conflate() and buffer(1, DROP_OLDEST)?
    conflate() is effectively a CONFLATED channel that keeps only the most recent value, dropping intermediates. buffer(1, DROP_OLDEST) behaves equivalently in practice — conflate() is the named shorthand for that latest-wins semantics.

saying these in an interview costs you the question

  • Saying channelFlow defaults to an unlimited buffer
  • Believing send() never suspends
  • Not knowing buffer() configures the backing channel
  • Confusing DROP_OLDEST with DROP_LATEST semantics
  • Unaware that adjacent buffer/flowOn operators fuse

context