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?
answer
- Reuse suspend instead of a second callback model
- Wins: sequential reads, free cancellation/context, auto backpressure
- Cost: cold = unicast ⇒ shareIn/stateIn to multicast
- No request(n); use buffer/conflate for capacity
- Interop via asFlow/asPublisher/asFlux adapters
basics
~20 sBuilding 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 sKotlin 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// 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
Can say Flow is built on coroutines so async code reads sequentially.
Lists concrete wins (auto backpressure, free cancellation) and knows shareIn makes a flow shared.
Articulates unicast-by-default vs hot multicast, no native request(n), and names the interop adapters.
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