skip to content

Why did Kotlin model Flow as cold and suspend-based rather than adopting the Reactive Streams Publisher/Subscriber callback model? What trade-offs does that create, including interop?

level: principalimportance: nice to knowfreq 30%

answer

  1. Reuse suspend instead of a second callback model
  2. Wins: sequential reads, free cancellation/context, auto backpressure
  3. Cost: cold = unicast ⇒ shareIn/stateIn to multicast
  4. No request(n); use buffer/conflate for capacity
  5. Interop via asFlow/asPublisher/asFlux adapters

basics

~20 s

Building Flow on coroutines lets async stream code read like normal sequential code, with built-in cancellation and backpressure via suspension. The trade-off is that the default Flow lacks multicasting and needs adapters to talk to RxJava/Reactor.

solid answer

~40 s

Kotlin already had coroutines and structured concurrency, so modeling Flow on suspend reuses one async primitive instead of adding a parallel callback world. Benefits: operators are plain suspend functions so pipelines read sequentially; backpressure is automatic via suspension; cancellation and context come free from structured concurrency; no Scheduler zoo. Trade-offs: a cold Flow is unicast (re-runs per collector) so multicasting needs shareIn/stateIn or hot SharedFlow/StateFlow; there's no spec-level request(n) demand signaling for fine-grained pull; and crossing into the Reactive Streams ecosystem needs the kotlinx-coroutines-reactive/reactor/rx adapters (asFlow/asPublisher/asFlux) which bridge demand-based and suspension-based backpressure. The design optimizes for readability and correctness within the coroutine world, accepting interop seams at the boundary.

code

kotlin · 10 lines
kotlin
// Cold (unicast) -> hot (multicast) to share one upstream:
val shared: SharedFlow<Quote> = quoteFlow
    .shareIn(
        scope = appScope,
        started = SharingStarted.WhileSubscribed(5_000),
        replay = 1
    )

// Interop: consume an RxJava/Reactor Publisher as a Flow
val asKotlin: Flow<Row> = reactiveRepository.findAll().asFlow()

go deeper

for a junior

Can say Flow is built on coroutines so async code reads sequentially.

for a middle

Lists concrete wins (auto backpressure, free cancellation) and knows shareIn makes a flow shared.

for a senior

Articulates unicast-by-default vs hot multicast, no native request(n), and names the interop adapters.

for a principal

Weighs the trade-offs holistically — readability/correctness vs multicasting/demand signaling — and reasons about SharingStarted policies, leak avoidance, and bridging two backpressure philosophies at ecosystem boundaries.

## The design choice Reactive Streams (RxJava, Reactor, the `Publisher`/`Subscriber`/`Subscription` spec) coordinate async streams with **callbacks** and an explicit **`request(n)` demand protocol**. Kotlin instead built `Flow` on **`suspend` + coroutines + structured concurrency**. The motivation: Kotlin *already* has suspension as its async primitive, so a stream type can reuse it instead of introducing a second, callback-shaped concurrency model. ## What the suspend/cold model buys - **Sequential-looking pipelines.** Operators (`map`, `filter`, `transform`) are ordinary functions over `suspend` lambdas, so async code reads top-to-bottom without callback nesting. - **Backpressure for free.** Suspension at `emit` paces the producer to the collector — no demand bookkeeping. - **Cancellation & context for free.** Structured concurrency means a collected flow inherits the collector's scope: cancel the scope, the flow stops; no `Disposable` to manage manually. - **No Scheduler proliferation.** Context is a `CoroutineContext`/dispatcher, changed with `flowOn`, instead of `subscribeOn`/`observeOn` Schedulers. - **Coldness = composability.** Each collection is a fresh, restartable run, which makes `retry`, re-collection, and per-subscriber parameterization natural. ## The trade-offs - **Unicast by default.** A cold `Flow` re-executes per collector; it does **not** multicast. To share one upstream among many collectors you convert to hot: `shareIn(scope, started, replay)` / `stateIn(...)`, or use `SharedFlow`/`StateFlow` directly. That's an explicit decision, not the default. - **No spec-level `request(n)`.** Fine-grained pull-based demand isn't surfaced; you express capacity with `buffer`/`conflate` instead. For most app code this is simpler; for strict Reactive Streams demand semantics it's a difference. - **Hot streams need lifecycle care.** `SharedFlow`/`StateFlow` are active regardless of collectors, so a `shareIn`/`stateIn` requires a scope and a `SharingStarted` policy (`Eagerly`/`Lazily`/`WhileSubscribed`) to avoid leaks. - **Interop seams.** To talk to the Reactive Streams ecosystem you use adapters from `kotlinx-coroutines-reactive`/`-reactor`/`-rx2`/`-rx3`: ```kotlin import kotlinx.coroutines.reactive.asFlow import kotlinx.coroutines.reactor.asFlux val flow: Flow<T> = somePublisher.asFlow() // Publisher -> Flow val flux: Flux<T> = someFlow.asFlux() // Flow -> Flux (needs a scope/context) ``` These adapters bridge two backpressure philosophies: the Reactive side's `request(n)` demand and Flow's suspension. The bridge is well-defined but is still a boundary you must cross deliberately. ## Net assessment The cold/suspend model **optimizes for readability, correctness, and unified concurrency** inside the Kotlin/coroutines world, at the cost of **built-in multicasting and native demand signaling**, both recoverable via hot flows and adapters. For application-level code (UI, services) this is usually the better default; for library authors interoperating with reactive infrastructures, the adapters and hot-flow tooling are where the real design work lives.

  • If a cold Flow is unicast, how do you share one expensive upstream across many collectors?
    Convert it to a hot flow with shareIn (or stateIn for state), supplying a scope, replay count, and a SharingStarted policy.
  • What does SharingStarted.WhileSubscribed buy you over Eagerly?
    It starts the upstream only while there are subscribers and stops it (after an optional timeout) when none remain, avoiding work/leaks when nobody is listening.
  • When bridging a Reactive Streams Publisher to Flow, what semantic gap must the adapter handle?
    It must translate the Publisher's request(n) demand-based backpressure into Flow's suspension-based backpressure.

Choosing suspend over callbacks is like writing prose instead of a chain of sticky notes: the story reads in order, but to broadcast it to a crowd (multicast) or hand it to a different publishing house (Rx) you need an explicit megaphone or translator.

saying these in an interview costs you the question

  • Claiming Flow implements the Reactive Streams request(n) spec natively
  • Saying a cold Flow multicasts to multiple collectors by default
  • Asserting Flow cannot interoperate with RxJava/Reactor at all
  • Ignoring the lifecycle/leak risk of hot flows from shareIn/stateIn
  • Treating buffer/conflate and request(n) as identical mechanisms

context