"Flowing Data" — Flow, Channels and Virtual Threads
Store managers want live stock numbers, and Kabir reaches for Flux. We cover cold Flow and its operators, back-pressure with buffer, conflate and debounce, hot StateFlow and SharedFlow, channels, Reactor interop in both directions, testing with runTest and virtual time, and a decision table for virtual threads versus coroutines.
Story Opening
The live stock board had been on the staging screen in the ops room for two days when someone noticed that the Mumbai tile hadn’t moved since breakfast. Nothing in the logs, no alert, and the Pune tiles were still ticking.
Kabir had written the feed the afternoon the request arrived. Lena’s remark that a Kotlin flow “reads like a loop” had sounded like Android talk, and the deadline was Friday, so he wrote it the way he had written a dozen WebFlux endpoints:
// Fragment of story/ReactorDashboard.kt// Kabir's first instinct, written the WebFlux way. Every step is an operator.fun liveStock(store: String): Flux<String> = Flux.interval(Duration.ofMillis(10)) .take(3) .flatMap { tick -> fetchLevel(store, tick) } .map { "${it.store} ${it.sku}: ${it.units}" } .onErrorResume { Flux.empty() } // "a bad reading shouldn't take the dashboard down"Replaying the gateway’s recorded traffic reproduced it in one run. Mumbai’s stream emitted a single reading and then completed. Its second payload was "12O", with a letter O; toInt() threw, onErrorResume turned the exception into an empty Flux, and every subscriber saw a stream that had ended successfully. He had added that line for resilience, as he had on his last three reactive projects.
Lena read the chain and asked a different question. “Which of these lines needs a reactive type, and which ones just need to wait?” Only the timer did. The rest was a loop with pauses in it, and in a loop an exception doesn’t disappear unless somebody writes the catch.
Java → Kotlin: The Quick Map
| Java / Reactor | Kotlin | Note |
|---|---|---|
Flux<T> (cold) | Flow<T> | Built on suspend; no Publisher protocol |
Mono<T> | suspend fun …(): T | One value doesn’t need a stream type |
Flux.generate(…) / a computed sequence | flow { emit(…) } | Sequential code with emit |
Flux.create(sink -> listener …) | callbackFlow { … awaitClose { } } | Bridges a Java listener API |
flatMap(x -> asyncCall(x)) | map { suspendCall(it) }, or flatMapMerge | map is sequential and ordered; flatMapMerge is concurrent (16 by default) |
concatMap | flatMapConcat | Opt-in (ExperimentalCoroutinesApi), like flatMapMerge |
onBackpressureBuffer / onBackpressureLatest | buffer() / conflate() | Back-pressure is suspension |
subscribeOn | flowOn(dispatcher) | Positional: moves only the operators above it |
publishOn | (no operator) | Downstream runs in the collector’s context; collect somewhere else to move it |
Sinks.many().replay().latestOrDefault(x) | MutableStateFlow(x) | Always has a value; also skips equal updates; update { } like updateAndGet |
Sinks.many().multicast().directBestEffort() | MutableSharedFlow | Broadcast; replay and buffer configurable |
BlockingQueue between threads | Channel | Each element to one receiver; suspends instead of blocking |
StepVerifier + virtual time | runTest + advanceTimeBy | |
Schedulers.boundedElastic() | Dispatchers.IO, or a virtual-thread dispatcher |
Conceptual Deep-Dive
A Flow is a suspend function that returns several times
Reactor builds streams out of a protocol: Publisher, Subscriber, request(n). Every operator is an object that implements that protocol, which is why the stack traces are deep and the code is a chain of operators. Kotlin’s Flow is much smaller. Its core is one interface with one function:
interface Flow<T> { suspend fun collect(collector: FlowCollector<T>) }interface FlowCollector<T> { suspend fun emit(value: T) }A flow { } block is a suspend function that calls emit several times, and collect runs it, passing in the code that handles each value. emit calls straight into that code, so a producer can’t get ahead of a slow collector: it is still inside emit until the collector returns. That is back-pressure: no request(n), no buffers, unless you add them. Operators such as map and filter are short functions that wrap one collector in another; the Technical Explanation writes one in three lines.
The shift: in Reactor, you describe a pipeline and hand it to a framework; in Kotlin, a flow is ordinary sequential code that pauses. Loops, try/catch, if and local variables all work inside it. You still get operators for concurrency and time (buffer, debounce, flatMapMerge), but you reach for them when you need them, not to write a loop.
Cold, hot, and queues
| Type | Who produces | Who receives | Value when nobody listens | Java/Reactor analogue |
|---|---|---|---|---|
Flow (cold) | The collector, by collecting | That collector only | Nothing runs | Cold Flux |
SharedFlow (hot) | Code that calls emit | Every current subscriber | Dropped, unless replay | Sinks.many().multicast().directBestEffort() |
StateFlow (hot) | Code that sets value | Every subscriber, latest value first | Kept: there is always a current value | Sinks.many().replay().latestOrDefault(x) |
Channel | send | Exactly one receiver per element | Buffered or suspended | BlockingQueue |
A useful rule: describe data with Flow, share state with StateFlow, broadcast events with SharedFlow, and hand work out with Channel.
Technical Explanation
Cold flows and operators
import kotlinx.coroutines.Dispatchersimport kotlinx.coroutines.delayimport kotlinx.coroutines.flow.Flowimport kotlinx.coroutines.flow.filterimport kotlinx.coroutines.flow.flowimport kotlinx.coroutines.flow.flowOnimport kotlinx.coroutines.flow.mapimport kotlinx.coroutines.flow.takeimport kotlinx.coroutines.flow.toListimport kotlinx.coroutines.runBlockingimport kotlinx.coroutines.withContext
data class StockLevel(val store: String, val sku: String, val units: Int)
// A cold flow: the block runs once per collector, when it collects. Nothing runs before that.fun stockLevels(store: String): Flow<StockLevel> = flow { println("scanner connected for $store") for (tick in 0..3) { delay(10) // suspends between values; no thread is held emit(StockLevel(store, "SHW-1001", 120 - tick * 7)) }}
fun main() = runBlocking { val pune = stockLevels("PUN-014") // nothing printed yet: building a flow runs nothing
// Operators read like the Stream API, but they are suspend-aware and run in the collector's coroutine. pune .filter { it.units < 115 } .map { "${it.store}: ${it.units}" } .take(2) .collect { println(it) } // -> scanner connected for PUN-014 // -> PUN-014: 113 // -> PUN-014: 106
// A second collection runs the producer again. println(pune.toList().size) // -> scanner connected for PUN-014 // -> 4
// flowOn changes where the UPSTREAM runs; the collector stays where it is. stockLevels("MUM-002") .map { Thread.currentThread().name.startsWith("DefaultDispatcher") } .flowOn(Dispatchers.Default) .take(1) .collect { println("produced on Default: $it") } // -> scanner connected for MUM-002 // -> produced on Default: true
// Emitting from a different context inside flow { } breaks the flow's contract. try { flow { withContext(Dispatchers.Default) { emit(1) } }.collect { } } catch (e: IllegalStateException) { println(e.message!!.substringBefore(":")) // -> Flow invariant is violated }}Building pune printed nothing: a flow is a recipe. Each collect runs the block again, which is why “scanner connected” appears twice. Operators run in the collector’s coroutine, in order, one element at a time, so take(2) cancels the producer after the second value, as first() stopped the sequence in Part 7.
flowOn(Dispatchers.Default) moves everything above it onto another dispatcher and leaves the collector where it was; the next section shows how. The last example shows the rule that makes flowOn necessary: a flow { } block must emit from the coroutine that collects it, and switching context inside it with withContext fails at runtime with “Flow invariant is violated”. Use flowOn instead.
Under the hood: operators are collectors, buffer is a channel
import kotlinx.coroutines.flow.Flowimport kotlinx.coroutines.flow.bufferimport kotlinx.coroutines.flow.flowimport kotlinx.coroutines.runBlocking
// An operator is a new flow whose collect collects the upstream. The library's map and filter// have the same shape; they use an internal builder (unsafeTransform) that skips the safety// checks flow { } performs on every emit, which this hand-written version keeps.fun <T, R> Flow<T>.mapped(transform: suspend (T) -> R): Flow<R> = flow { collect { value -> emit(transform(value)) }}
fun readings() = flow { println("emit 1"); emit(1) println("emit 2"); emit(2)}
fun main() = runBlocking { // emit calls straight down into the collector: each value goes all the way through // before the producer gets the next line. One coroutine, one call stack. readings().mapped { println("map $it"); it * 10 }.collect { println("collect $it") } // -> emit 1 // -> map 1 // -> collect 10 // -> emit 2 // -> map 2 // -> collect 20
// buffer splits the flow in two: everything above it runs in a new producer coroutine, // which sends into a channel; the collector receives from it. Now the producer runs ahead. readings().mapped { println("map $it"); it * 10 }.buffer().collect { println("collect $it") } // -> emit 1 // -> map 1 // -> emit 2 // -> map 2 // -> collect 10 // -> collect 20}mapped is the whole operator model, minus the shortcuts the library takes for speed (see the comment). It returns a new flow; collecting that flow collects the upstream, transforms each value and emits it downstream. Nothing is queued. emit is a plain suspend call into the next collector, so without a buffer the first trace interleaves strictly and a slow consumer slows the producer to its pace. A Reactor operator is a Subscriber plus a Subscription with request accounting; a Flow operator is a function calling a function.
buffer() changes that, and so does every operator that needs two things to happen at once: conflate, debounce, flatMapMerge, and flowOn when it actually changes the dispatcher. Each one splits the pipeline into two coroutines joined by a channel:
The second trace shows the split: both emits and both maps run before the first collect. The channel’s capacity is now the back-pressure, because send suspends the producer when the channel is full. flowOn is the same split with a different dispatcher on the producer side, which is why it affects only the operators above it. Reactor’s publishOn, which moves everything below it, has no operator twin: to move the collector, you collect somewhere else.
The opening’s Flux, as a Flow
import kotlinx.coroutines.ExperimentalCoroutinesApiimport kotlinx.coroutines.delayimport kotlinx.coroutines.flow.Flowimport kotlinx.coroutines.flow.flatMapConcatimport kotlinx.coroutines.flow.flatMapMergeimport kotlinx.coroutines.flow.flowimport kotlinx.coroutines.flow.flowOfimport kotlinx.coroutines.flow.mapimport kotlinx.coroutines.flow.toListimport kotlinx.coroutines.test.runTest
// The gateway client as a suspend function: one value, so no Mono. (StockLevel is from ColdFlows.kt.)suspend fun fetchLevel(store: String, tick: Int): StockLevel { delay(5) // the remote call val raw = if (store == "MUM-002" && tick == 1) "12O" else "${120 - tick * 7}" return StockLevel(store, "SHW-1001", raw.toInt())}
// The opening's Flux chain, rewritten. Only the timer needed a stream; the rest is a loop that waits.fun liveStock(store: String): Flow<String> = flow { for (tick in 0 until 3) { delay(10) emit(fetchLevel(store, tick)) // Reactor needed flatMap for this call }}.map { "${it.store} ${it.sku}: ${it.units}" }
// A lookup whose latency depends on the item: later items answer first.fun slowLookup(n: Int): Flow<Int> = flow { delay((4 - n) * 10L) emit(n)}
@OptIn(ExperimentalCoroutinesApi::class) // flatMapMerge and flatMapConcat need opt-in in kotlinx.coroutines 1.11.0fun main() = runTest { liveStock("PUN-014").collect { println(it) } // -> PUN-014 SHW-1001: 120 // -> PUN-014 SHW-1001: 113 // -> PUN-014 SHW-1001: 106
// No operator swallows anything by default: the failure reaches the code that collected. try { liveStock("MUM-002").collect { println(it) } } catch (e: NumberFormatException) { println("MUM-002 failed: ${e.message}") } // -> MUM-002 SHW-1001: 120 // -> MUM-002 failed: For input string: "12O"
// map with a suspend call: one at a time, in order. This is what most Reactor flatMaps meant. println(flowOf(1, 2, 3).map { n -> slowLookup(n).toList().single() }.toList()) // -> [1, 2, 3] // flatMapMerge: concurrent, so results arrive in completion order, like Reactor's flatMap. println(flowOf(1, 2, 3).flatMapMerge { slowLookup(it) }.toList()) // -> [3, 2, 1] // flatMapConcat: one inner flow after another, like Reactor's concatMap. println(flowOf(1, 2, 3).flatMapConcat { slowLookup(it) }.toList()) // -> [1, 2, 3]}That answers Lena’s question. The timer became a for loop with a delay, and the remote call is a suspend function called inside it. Reactor needed flatMap for that call because a Mono has to be subscribed to; a suspend function is just called.
The Mumbai failure now does what an exception in a loop does: it ends the flow and reaches the code that collected it, with its message intact. No operator in kotlinx.coroutines swallows an exception by default. If the dashboard should show a fallback instead, that is a catch { } operator someone writes and a reviewer can see; it handles exceptions from upstream only, and the Gotchas show why that matters.
The flatMap lines are the other half of the habit. Reactor’s flatMap subscribes to the inner publishers concurrently, so results arrive in completion order. In a Flow, map with a suspend call does the per-element call one at a time, in order. Neither is free: one reorders, the other serialises (see the Gotchas before porting a chain). Reach for flatMapMerge when you need the concurrency back, and expect completion order ([3, 2, 1]); flatMapConcat is Reactor’s concatMap.
Back-pressure, time and virtual time
import kotlinx.coroutines.ExperimentalCoroutinesApiimport kotlinx.coroutines.FlowPreviewimport kotlinx.coroutines.delayimport kotlinx.coroutines.flow.bufferimport kotlinx.coroutines.flow.collectLatestimport kotlinx.coroutines.flow.conflateimport kotlinx.coroutines.flow.debounceimport kotlinx.coroutines.flow.flowimport kotlinx.coroutines.flow.toListimport kotlinx.coroutines.test.TestScopeimport kotlinx.coroutines.test.currentTimeimport kotlinx.coroutines.test.runTest
// A scanner that emits a reading every 100 ms (of virtual time: runTest skips the waiting).// Each consumer below takes 300 ms per reading.fun scans(count: Int) = flow { repeat(count) { delay(100) emit(it) }}
// Runs a block and prints what it collected and how much virtual time it took.@OptIn(ExperimentalCoroutinesApi::class) // currentTime is experimental in kotlinx-coroutines-test 1.11.0suspend fun TestScope.timed(label: String, block: suspend () -> List<Int>) { val start = currentTime val seen = block() println("$label: saw $seen in ${currentTime - start} ms")}
@OptIn(FlowPreview::class) // debounce is a Flow preview API in kotlinx.coroutines 1.11.0fun main() = runTest { // No buffer: producer and consumer take turns, 100 + 300 ms per reading. timed("sequential") { buildList { scans(5).collect { delay(300); add(it) } } } // -> sequential: saw [0, 1, 2, 3, 4] in 2000 ms
// buffer: the producer runs ahead in its own coroutine; the consumer sets the pace. timed("buffer") { buildList { scans(5).buffer().collect { delay(300); add(it) } } } // -> buffer: saw [0, 1, 2, 3, 4] in 1600 ms
// conflate: the consumer only ever sees the latest value; intermediate ones are dropped. timed("conflate") { buildList { scans(5).conflate().collect { delay(300); add(it) } } } // -> conflate: saw [0, 2, 4] in 1000 ms
// collectLatest: a new value cancels the work on the previous one. timed("collectLatest") { buildList { scans(5).collectLatest { delay(300); add(it) } } } // -> collectLatest: saw [4] in 800 ms
// debounce: emit only after 250 ms of quiet. A burst of scans becomes one update. val burst = flow { emit(1); delay(50); emit(2); delay(50); emit(3) // a burst delay(400) emit(4) // a lone scan } timed("debounce") { burst.debounce(250).toList() } // -> debounce: saw [3, 4] in 500 ms // (the last value goes out as soon as the flow completes, without waiting 250 ms)}This example runs under runTest, from kotlinx-coroutines-test, which uses a virtual clock: delay(300) advances the clock without waiting, so the whole file runs in milliseconds and the timings are exact. (Several examples in this part use runTest in main for that reason alone; it is a test dependency, and production code uses a real scope.) With a producer emitting every 100 ms and a consumer taking 300 ms:
- No buffer: producer and consumer take turns, 400 ms per value.
buffer(): the producer runs ahead into the channel, and the consumer’s pace decides the total.conflate(): a channel with one slot that the producer overwrites, so the consumer gets only the latest value each time it’s ready. Values 1 and 3 are dropped. That’s right for a dashboard, where only the current stock matters.collectLatest: a new value cancels the processing of the previous one. Only the last value completes.debounce(250): a value is emitted only after 250 ms without a newer one, which turns a burst of scans into one update. When the flow completes, the last value goes out straight away.
Each example opts in where kotlinx.coroutines 1.11.0 requires it (@OptIn(ExperimentalCoroutinesApi::class) or @OptIn(FlowPreview::class)); buffer, conflate, collectLatest, callbackFlow and runTest itself are stable.
Hot flows: StateFlow and SharedFlow
import kotlinx.coroutines.delayimport kotlinx.coroutines.flow.MutableSharedFlowimport kotlinx.coroutines.flow.MutableStateFlowimport kotlinx.coroutines.flow.SharingStartedimport kotlinx.coroutines.flow.asStateFlowimport kotlinx.coroutines.flow.flowimport kotlinx.coroutines.flow.stateInimport kotlinx.coroutines.flow.takeimport kotlinx.coroutines.flow.toListimport kotlinx.coroutines.launchimport kotlinx.coroutines.test.runTest
class ShelfSensor { // StateFlow: always has a value; new collectors get the current one; equal updates are skipped. private val _units = MutableStateFlow(120) val units = _units.asStateFlow() // read-only view for the outside world
fun scan(newUnits: Int) { _units.value = newUnits }}
fun main() = runTest { val sensor = ShelfSensor() val seen = mutableListOf<Int>() val watcher = launch { sensor.units.toList(seen) } delay(1) for (units in listOf(118, 118, 110, 110, 95)) { sensor.scan(units) delay(1) // give the watcher time to see every update, so nothing is conflated } watcher.cancel() println(seen) // -> [120, 118, 110, 95] // (the repeated 118 and 110 were equal to the current value, so they emitted nothing)
// Mutating the object inside a StateFlow emits nothing: the reference didn't change. val skus = MutableStateFlow(mutableListOf("SHW-1001")) var updates = 0 val counter = launch { skus.collect { updates++ } } delay(1) skus.value.add("SHW-2040") // invisible to collectors delay(1) counter.cancel() println("updates seen: $updates") // -> updates seen: 1
// SharedFlow: a broadcast of events. No current value; replay decides what late subscribers get. val scans = MutableSharedFlow<String>(replay = 1) scans.emit("PUN-014 aisle 7") scans.emit("PUN-014 aisle 8") val late = scans.take(1).toList() println(late) // -> [PUN-014 aisle 8]
// tryEmit (the non-suspending emit) on a SharedFlow with no buffer: with nobody subscribed, // the value is dropped and tryEmit reports success; with a subscriber, it fails. val alerts = MutableSharedFlow<String>() println(alerts.tryEmit("low stock: SHW-2040")) // -> true val listener = launch { alerts.collect { } } delay(1) println(alerts.tryEmit("low stock: SHW-3001")) // -> false listener.cancel()
// stateIn: one upstream shared by every collector. WhileSubscribed(5_000) starts it with the // first subscriber and stops it 5 s (of virtual time, here) after the last one leaves. var polls = 0 val stockPoll = flow { polls++ // one database poll loop per start while (true) { emit(100) delay(1_000) } } val shared = stockPoll.stateIn(backgroundScope, SharingStarted.WhileSubscribed(5_000), initialValue = 0) val dashboards = List(2) { launch { shared.collect { } } } delay(10) println("dashboards=2, polls=$polls") // -> dashboards=2, polls=1 dashboards.forEach { it.cancel() } delay(6_000) // nobody subscribed for more than 5 s: the poll stops val again = launch { shared.collect { } } delay(10) println("after a quiet spell, polls=$polls") // -> after a quiet spell, polls=2 again.cancel()}StateFlow is an observable value, and the workhorse of a dashboard like this one. It always has a current value, new collectors start with it (120), and setting a value equal to the current one emits nothing: the repeated 118 and 110 never reached the watcher, although it had time to see each update. Mutating the object inside it emits nothing either, for a different reason (see the Gotchas). Expose it read-only with asStateFlow(), the same backing-property pattern as Part 3. Change it with update { }, which retries a compare-and-set until it wins, like AtomicReference.updateAndGet; value = value + x is a race between concurrent writers. SharedFlow is for events that matter even when they repeat. Without replay, a value emitted while nobody is subscribed is gone; with replay = 1, a late subscriber gets the last one.
stateIn and shareIn convert a cold flow into a hot one shared by every collector, so an expensive upstream such as a database poll runs once however many dashboards are open (polls=1 for two). stateIn keeps the latest value and takes an initial one, or, in its suspending overload stateIn(scope), waits for the first; shareIn has no current value, and its replay decides what late subscribers get. SharingStarted.Eagerly starts at once; WhileSubscribed(5_000) starts with the first subscriber and stops 5 seconds after the last one leaves, which is why the quiet spell above cost a second poll start. The grace period rides out a browser refresh without restarting the upstream. Both need a scope that outlives the collectors: in tests, backgroundScope; in Spring, one the bean owns (see the Tips).
Channels
import kotlinx.coroutines.ExperimentalCoroutinesApiimport kotlinx.coroutines.channels.Channelimport kotlinx.coroutines.channels.produceimport kotlinx.coroutines.coroutineScopeimport kotlinx.coroutines.delayimport kotlinx.coroutines.launchimport kotlinx.coroutines.runBlockingimport java.util.concurrent.ConcurrentHashMap
// A channel is a queue between coroutines: each element goes to exactly ONE receiver.// That makes it the tool for distributing work, where a Flow is the tool for describing a stream.@OptIn(ExperimentalCoroutinesApi::class) // produce { } is experimental in kotlinx.coroutines 1.11.0fun main() = runBlocking { val labelsPrinted = ConcurrentHashMap<String, Int>()
coroutineScope { val jobs = produce(capacity = Channel.BUFFERED) { for (aisle in 1..12) send("aisle-$aisle") } // closed automatically when the producer finishes
repeat(3) { printer -> launch { for (job in jobs) { // fan-out: three printers share the queue delay(10) labelsPrinted.merge("printer-$printer", 1, Int::plus) } } } } println("${labelsPrinted.values.sum()} labels, ${labelsPrinted.size} printers busy") // -> 12 labels, 3 printers busy}A Channel is a coroutine BlockingQueue: send suspends when the buffer is full, receive suspends when it is empty, and each element goes to exactly one receiver. produce { } creates a producer coroutine and a channel that closes when the producer finishes, so the printers’ for loops end by themselves. Use channels to distribute work between coroutines, as here, or to pass messages to a single coroutine that owns some state. Use flows for everything that is “a stream of values to transform”. Most application code meets channels only inside builders such as the next one.
Java listener APIs as flows: callbackFlow
Plenty of Java SDKs push data through listeners: a scanner SDK, a JMS MessageListener, a WebSocket client. The Reactor answer is Flux.create(sink -> …); Kotlin’s is callbackFlow, a flow builder with a channel inside:
import java.util.List;import java.util.concurrent.CopyOnWriteArrayList;
// A vendor's shelf-scanner SDK: listener-based, and the caller must unregister by hand.public final class ScannerSdk { public interface ScanListener { void onScan(String sku, int units);
void onError(Exception e); }
private final List<ScanListener> listeners = new CopyOnWriteArrayList<>();
public void addListener(ScanListener listener) { listeners.add(listener); }
public void removeListener(ScanListener listener) { listeners.remove(listener); }
public int listenerCount() { return listeners.size(); }
// The real SDK calls listeners from its own I/O thread; the example calls this directly. public void deliver(String sku, int units) { for (ScanListener listener : listeners) listener.onScan(sku, units); }}import kotlinx.coroutines.asyncimport kotlinx.coroutines.channels.awaitCloseimport kotlinx.coroutines.flow.Flowimport kotlinx.coroutines.flow.callbackFlowimport kotlinx.coroutines.flow.takeimport kotlinx.coroutines.flow.toListimport kotlinx.coroutines.runBlockingimport kotlinx.coroutines.yield
data class Scan(val sku: String, val units: Int)
// A Java listener API as a Flow. callbackFlow gives the block a channel: callbacks push into it,// the collector reads from it, and awaitClose keeps the listener registered until the collector stops.fun ScannerSdk.scans(): Flow<Scan> = callbackFlow { val listener = object : ScannerSdk.ScanListener { override fun onScan(sku: String, units: Int) { trySend(Scan(sku, units)) // callbacks can't suspend, so the non-suspending send }
override fun onError(e: Exception) { close(e) // fails the flow with the SDK's exception } } addListener(listener) awaitClose { removeListener(listener) } // runs on completion, failure or cancellation}
fun main() = runBlocking { val sdk = ScannerSdk() val firstTwo = async { sdk.scans().take(2).toList() } while (sdk.listenerCount() == 0) yield() // let the collector start and register
sdk.deliver("SHW-1001", 40) sdk.deliver("SHW-2040", 8) sdk.deliver("SHW-3001", 75) // buffered, then discarded when take(2) cancels the flow
println(firstTwo.await()) // -> [Scan(sku=SHW-1001, units=40), Scan(sku=SHW-2040, units=8)] println("listeners after: ${sdk.listenerCount()}") // -> listeners after: 0}The listener pushes with trySend, the non-suspending send, because a Java callback can’t suspend. If the channel’s buffer (64 elements by default) is full, trySend fails and the value is dropped; choose the policy explicitly with .buffer(capacity) or .conflate() on the resulting flow. close(e) fails the flow, and the collector sees the SDK’s exception. awaitClose { } suspends the builder until the flow ends, for whatever reason, and then runs the cleanup: here, take(2) cancelled the flow and the listener was removed (listeners after: 0). Forget awaitClose and the flow fails as soon as the block returns, with “‘awaitClose { yourCallbackOrListener.cancel() }’ should be used in the end of callbackFlow block”, which is kinder than a listener that leaks for the life of the process.
Reactor interop
import kotlinx.coroutines.delayimport kotlinx.coroutines.flow.flowimport kotlinx.coroutines.flow.mapimport kotlinx.coroutines.flow.toListimport kotlinx.coroutines.reactive.asFlowimport kotlinx.coroutines.reactive.awaitSingleimport kotlinx.coroutines.reactor.awaitSingleOrNullimport kotlinx.coroutines.reactor.asFluximport kotlinx.coroutines.reactor.fluximport kotlinx.coroutines.reactor.monoimport kotlinx.coroutines.runBlockingimport reactor.core.publisher.Fluximport reactor.core.publisher.Mono
// A Reactor-based client from another team.fun legacyPrice(sku: String): Mono<Long> = Mono.just(if (sku == "SHW-1001") 16_500L else 72_000L)
fun legacyStock(): Flux<Int> = Flux.just(120, 113, 106)
fun main() = runBlocking { // Reactor -> coroutines: await a Mono, collect a Flux as a Flow. println(legacyPrice("SHW-1001").awaitSingle()) // -> 16500 // An empty Mono ("not found" in most Java code) makes awaitSingle throw; use awaitSingleOrNull. println(Mono.empty<Long>().awaitSingleOrNull()) // -> null try { Mono.empty<Long>().awaitSingle() } catch (e: NoSuchElementException) { println("awaitSingle on empty: ${e::class.simpleName}") // -> awaitSingle on empty: NoSuchElementException } println(legacyStock().asFlow().map { it - 100 }.toList()) // -> [20, 13, 6]
// Coroutines -> Reactor: suspend code behind a Mono or a Flux, for APIs that expect them. val price: Mono<String> = mono { "₹${legacyPrice("SHW-2040").awaitSingle() / 100}" } println(price.block()) // -> ₹720
val ticks: Flux<Int> = flux { for (i in 1..3) { delay(5) send(i) // flux { } is a producer: send, not emit } } println(ticks.collectList().block()) // -> [1, 2, 3]
// A Flow as a Flux, e.g. to return from a WebFlux controller that expects Publisher types. println(flow { emit("a"); emit("b") }.asFlux().collectList().block()) // -> [a, b]}kotlinx-coroutines-reactor connects the two models in both directions. awaitSingle() suspends until a Mono emits. Java code routinely returns Mono.empty() for “not found”, and awaitSingle() throws NoSuchElementException on it; awaitSingleOrNull() returns null, which ?: handles (Part 2). asFlow() collects a Flux as a Flow, with back-pressure mapped onto suspension. mono { } and flux { } build Reactor types from suspend code, and asFlux() exposes a Flow to Reactor. This is how Kotlin services coexist with Reactor-based libraries, and how Spring WebFlux runs suspend controller methods and Flow return types: it adapts them with these functions. Reactor’s Context, where WebFlux keeps tracing and security data, is visible to coroutines as a ReactorContext element; Part 14 wires up propagation for logging and tracing.
Virtual threads or coroutines?
Since Java 21, the JVM itself can make blocking cheap, which is the question Kabir has been holding since the virtual-thread demo that opened Part 10. The two work at different layers, and they combine:
import kotlinx.coroutines.asCoroutineDispatcherimport kotlinx.coroutines.asyncimport kotlinx.coroutines.awaitAllimport kotlinx.coroutines.runBlockingimport kotlinx.coroutines.withContextimport java.util.concurrent.Executors
// A blocking SDK call that can't be rewritten: it holds its thread for the whole call.fun blockingSupplierCall(i: Int): Int { Thread.sleep(100) return i}
// The two models combine: coroutines for structure, a virtual-thread dispatcher for blocking calls.fun main() = runBlocking { Executors.newVirtualThreadPerTaskExecutor().asCoroutineDispatcher().use { virtualThreads -> val results = withContext(virtualThreads) { (1..200).map { i -> async { blockingSupplierCall(i) } }.awaitAll() } println("${results.size} blocking calls, sum ${results.sum()}") // -> 200 blocking calls, sum 20100 println(withContext(virtualThreads) { Thread.currentThread().isVirtual }) // -> true }}asCoroutineDispatcher() turns a virtual-thread executor into a dispatcher. Coroutines give the 200 calls a scope, cancellation and awaitAll; virtual threads absorb the blocking SDK call without anyone choosing and sizing a thread pool for it. What virtual threads remove is the pool’s accidental limit, not the need for one: 200 unbounded calls to a supplier’s API is a rate-limit incident, so cap them with a Semaphore or limitedParallelism, as in Part 10.
| Situation | Choose | Why |
|---|---|---|
| Existing Spring MVC + JPA/JDBC service, blocking libraries throughout | Virtual threads (spring.threads.virtual.enabled=true) | Same code, cheaper waiting; nothing to rewrite. The connection pool still caps database concurrency |
| Service whose work is mostly streaming or many concurrent outbound calls, on WebFlux | Coroutines (suspend controllers, Flow) | Reads like blocking code on a non-blocking stack; Spring MVC also accepts suspend handlers |
| Fan-out with timeouts, cancellation, “first successful” or all-or-nothing | Coroutines | Structured concurrency is the point; Java’s StructuredTaskScope is still a preview API |
| Streams of values, live updates, back-pressure | Coroutines (Flow) | Virtual threads have no stream model |
| A blocking SDK inside coroutine code | Coroutines + a virtual-thread (or IO) dispatcher | Structure from coroutines, cheap blocking from the JVM |
| CPU-bound work | Neither helps; Dispatchers.Default or a bounded pool | Waiting isn’t the bottleneck |
| Java-heavy team, Kotlin only at the edges | Virtual threads | No colouring, no new concepts for Java callers |
Two caveats that apply in October 2026. Virtual threads in JDK 25 no longer pin their carrier thread inside synchronized blocks (fixed in JDK 24), but native calls and some class initialisation still pin; JFR’s jdk.VirtualThreadPinned event finds them. And Spring Boot 4.1 supports both models; Part 13 uses virtual threads for the MVC catalog service, and Part 14 uses coroutines for the reactive parts.
Step-by-Step Hands-On: A Live Stock Board
Code: kotlin-for-java-survivors/language/part11-flowing-data (file dashboard/StockBoard.kt, test StockBoardTest.kt).
Kabir builds the core of the store managers’ dashboard: shelf scans come in as events, stock levels go out as state, and low-stock alerts fire only when the set of low items changes.
Steps 1 to 4 — A scan, state, derived alerts, and a WebFlux view. The board holds an immutable map in a MutableStateFlow, exposed read-only. record uses update { }, which retries a compare-and-set until it wins, like AtomicReference.updateAndGet, so scans arriving from many coroutines at once can’t lose each other’s writes (value = value + x would race). Alerts are a derived flow: map to the set of low SKUs, and distinctUntilChanged so a scan that doesn’t change the set doesn’t produce an alert. The last line exposes the alerts as a Flux for Reactor-based callers; Part 14 serves a Flow like this one as server-sent events:
// Fragment of dashboard/StockBoard.kt// Step 1: a scan from a shelf scanner.data class ShelfScan(val store: String, val sku: String, val units: Int)
// Step 2: the board holds state: the latest units per "store/sku", as an immutable map.class StockBoard(private val reorderAt: Int = 10) { private val _levels = MutableStateFlow<Map<String, Int>>(emptyMap()) val levels: StateFlow<Map<String, Int>> = _levels.asStateFlow()
// Atomic read-modify-write: safe when scans arrive from many coroutines (or threads) at once. fun record(scan: ShelfScan) = _levels.update { it + ("${scan.store}/${scan.sku}" to scan.units) }
// Step 3: alerts derived from the state, emitted only when the low-stock set changes. val lowStock: Flow<Set<String>> = levels .map { all -> all.filterValues { it <= reorderAt }.keys } .distinctUntilChanged()}
// Step 4: the same stream as a Flux, for Reactor-based callers.fun StockBoard.lowStockFeed(): Flux<Set<String>> = lowStock.asFlux()Step 5 — Drive it.
// Fragment of dashboard/StockBoard.kt// Step 5: drive it under virtual time.@OptIn(ExperimentalCoroutinesApi::class) // runCurrent is experimental in kotlinx-coroutines-test 1.11.0fun main() = runTest { val board = StockBoard() val alerts = mutableListOf<Set<String>>() backgroundScope.launch { board.lowStock.toList(alerts) }
for (scan in listOf( ShelfScan("PUN-014", "SHW-1001", 40), ShelfScan("PUN-014", "SHW-2040", 8), ShelfScan("MUM-002", "SHW-1001", 5), ShelfScan("PUN-014", "SHW-2040", 60), )) { board.record(scan) runCurrent() // let the collector see each state; a StateFlow keeps only the latest }
println(board.levels.value) // -> {PUN-014/SHW-1001=40, PUN-014/SHW-2040=60, MUM-002/SHW-1001=5} alerts.forEach(::println) // -> [] // -> [PUN-014/SHW-2040] // -> [PUN-014/SHW-2040, MUM-002/SHW-1001] // -> [MUM-002/SHW-1001]}The first alert is the empty set, the initial state. PUN-014/SHW-2040 dropping to 8 adds it; MUM-002/SHW-1001 at 5 adds a second; the restock to 60 removes the first. The 40 scan never produced an alert, because the low-stock set didn’t change. The runCurrent() after each scan lets the collector run; without it, the StateFlow would conflate the four scans and the collector would see only the final state, which for a live dashboard is exactly right. If every intermediate value mattered, scans would go through a SharedFlow with a buffer instead.
Step 6 — Test it.
import com.shelfwise.part11.dashboard.ShelfScanimport com.shelfwise.part11.dashboard.StockBoardimport kotlinx.coroutines.ExperimentalCoroutinesApiimport kotlinx.coroutines.delayimport kotlinx.coroutines.flow.firstimport kotlinx.coroutines.launchimport kotlinx.coroutines.test.advanceTimeByimport kotlinx.coroutines.test.runCurrentimport kotlinx.coroutines.test.runTestimport org.junit.jupiter.api.Assertions.assertEqualsimport org.junit.jupiter.api.Test
@OptIn(ExperimentalCoroutinesApi::class) // advanceTimeBy, runCurrent, currentTimeclass StockBoardTest { @Test fun `a scan at or below the reorder point raises an alert`() = runTest { val board = StockBoard(reorderAt = 10)
board.record(ShelfScan("PUN-014", "SHW-2040", 10))
assertEquals(setOf("PUN-014/SHW-2040"), board.lowStock.first()) }
@Test fun `virtual time skips the waiting`() = runTest { val board = StockBoard() launch { delay(60_000) // a scanner that reports once a minute board.record(ShelfScan("MUM-002", "SHW-1001", 3)) }
advanceTimeBy(60_001) // instant in real time runCurrent()
assertEquals(3, board.levels.value["MUM-002/SHW-1001"]) assertEquals(60_001, testScheduler.currentTime) }}runTest runs a test body in a TestScope with a virtual clock. delay(60_000) inside it takes no real time, and advanceTimeBy and runCurrent move the clock and run whatever is due. Two rules keep these tests honest. Collectors of infinite flows, such as a StateFlow, belong in backgroundScope, which is cancelled when the test ends; in the test’s own scope, runTest would wait for them forever and fail with a timeout. And code under test must use the test’s dispatcher: a hard-coded Dispatchers.IO inside the class escapes virtual time. Inject a dispatcher or a scope instead. StockBoard launches nothing of its own, which is why it tests so easily.
Tips, Tricks & Gotchas
Tip — give a Spring bean its own scope for
stateInandshareIn. A fieldprivate val scope = CoroutineScope(SupervisorJob() + Dispatchers.Default)and a@PreDestroy fun close() = scope.cancel()tie the shared upstream to the bean’s lifecycle.GlobalScopewould outlive the application context, and a request’s scope would stop the upstream when the first request ends.
Tip —
runTestbelongs in tests. This part’s examples call it frommainonly to make timings deterministic. Production code collects in a real scope: a controller method, a bean scope, orrunBlockingat the edge of a command-line tool.
Gotcha — porting
flatMapchanges either the order or the latency. A Reactor.flatMap { fetch(it) }runs up to 256 inner calls at once. Rewritten asmap { fetch(it) }, the same chain makes one call at a time: correct, ordered, and as slow as the sum of the calls.flatMapMergerestores concurrency, but its default limit is 16 (the system propertykotlinx.coroutines.flow.defaultConcurrencychanges it), and results arrive in completion order. Decide which you need per chain, and passconcurrencyexplicitly when it matters. Reactor’sflatMapSequential(concurrent, but in source order) has no built-in equivalent.
Gotcha — mutating the object inside a
StateFlowemits nothing.skus.value.add("SHW-2040")changes the list but not the reference, and the equality check sees the same list (updates seen: 1, only the initial value). Java code that “updates the state” by mutating a shared collection compiles and silently stops updating the UI. Keep immutable values in aStateFlowand change them withupdate { }.
Gotcha —
tryEmitis not Reactor’stryEmitNext. Reactor’sSinks.many().multicast().directBestEffort()reportsFAIL_ZERO_SUBSCRIBERwhen nobody is listening. AMutableSharedFlowwithout a buffer does the opposite:tryEmitreturnstruewhen nobody is subscribed, because dropping the value succeeded, andfalseas soon as someone is, because delivering it would mean suspending (both lines inHotFlows.kt). Code ported from Reactor that treatstrueas “delivered” is wrong in both directions. Use the suspendingemit, or give the flow a buffer (extraBufferCapacity) and handlefalse.
Gotcha — a
try/catcharoundemitcatches the collector’s exceptions.emitruns the collector’s code. A Java-styletry { emit(x) } catch (e: Exception) { emit(fallback) }catches the consumer’s failure and tries to keep emitting, which fails with “Flow exception transparency is violated”. Use thecatch { }operator, which only sees upstream failures:
import kotlinx.coroutines.flow.catchimport kotlinx.coroutines.flow.flowimport kotlinx.coroutines.runBlocking
fun main() = runBlocking { // A Java-style try/catch around emit also catches the COLLECTOR's exceptions. // Emitting again after that breaks the flow's contract. val readings = flow { try { emit(1) } catch (e: IllegalStateException) { emit(-1) // "fallback" } } try { readings.collect { if (it == 1) throw IllegalStateException("display offline") } } catch (e: IllegalStateException) { println(e.message!!.lineSequence().first()) // -> Flow exception transparency is violated: }
// catch { } only sees upstream exceptions; the collector's own failure goes straight to the caller. try { flow { emit(1) } .catch { println("never reached") } .collect { throw IllegalStateException("display offline") } } catch (e: IllegalStateException) { println("collector failed: ${e.message}") // -> collector failed: display offline }}Gotcha — there are two
Flows on a Java classpath.java.util.concurrent.Flowis the JDK 9 holder class for the Reactive Streams interfaces (Flow.Publisher,Flow.Subscriber). When auto-import offers both, a Java developer’s fingers pick the familiar one. The result reads like a typo:fun levels(): Flow<Int>fails with “no type arguments expected for ‘class Flow : Any’”. You wantkotlinx.coroutines.flow.Flow; if you really have a JDKPublisher, kotlinx-coroutines-jdk9 converts it withasFlow().
Gotcha — a
Flowre-runs for every collector. A JavaStreamcan be consumed once; a coldFlowcan be collected any number of times, and each collection runs the producer again (“scanner connected” twice above). Two dashboards collectingstockLevels("PUN-014")open two scanner connections. Share it withstateInorshareIn.
Debugging and Troubleshooting
| Symptom | Likely cause | Fix |
|---|---|---|
| A stream stops updating, with no error anywhere | An error handler that completes the stream: Reactor’s onErrorResume { Flux.empty() }, or a catch { } that emits nothing | Log in the handler and emit a visible fallback, or let it fail and restart deliberately with retry/retryWhen |
IllegalStateException: Flow invariant is violated | emit called from a different coroutine context inside flow { }, usually via withContext | Use flowOn for the upstream, or channelFlow { } if you really need several producers |
| A collector never returns | Collecting a StateFlow or SharedFlow: they never complete | first { }, take(n), or collect in a scope you cancel |
| A test hangs, then fails with a timeout | An infinite collector launched in the runTest scope | Launch it in backgroundScope |
| Dashboard misses intermediate values | StateFlow and conflate() keep only the latest, by design | SharedFlow with a buffer, if every value matters |
| Virtual time doesn’t move in a test | The code under test uses its own Dispatchers.Default/IO | Inject the dispatcher or scope; pass the test’s |
NoSuchElementException from awaitSingle() | The Mono was empty | awaitSingleOrNull() |
| Warning: “this declaration needs opt-in” or “is in a preview state” | Experimental or preview coroutines API | @OptIn(ExperimentalCoroutinesApi::class) or @OptIn(FlowPreview::class), deliberately |
Key Takeaways
| Concept | Remember |
|---|---|
Flow | Cold, sequential, suspend-based; each collection runs the producer |
| Operators | A flow that collects its upstream; emit calls the collector directly |
| Back-pressure | Suspension at emit; buffer, conflate, collectLatest, debounce (preview) |
flatMap habits | map { suspendCall() } serialises; flatMapMerge reorders, 16 at a time by default; choose per chain |
buffer / flowOn | Split the pipeline into two coroutines and a channel; flowOn moves only what’s above it; never withContext around emit |
StateFlow | Observable value; skips equal values; immutable data changed with update { } |
SharedFlow | Hot event broadcast; replay decides what late subscribers see |
stateIn / shareIn | Share one upstream between many collectors, in a scope you own |
Channel | A suspending queue: each element to one receiver; for distributing work |
callbackFlow | Java listener → Flow: trySend in the callback, awaitClose { unregister } at the end |
| Reactor interop | awaitSingle, asFlow, mono { }, flux { }, asFlux |
| Testing | runTest virtual time; infinite collectors in backgroundScope; inject dispatchers |
| Virtual threads vs coroutines | Virtual threads for blocking stacks; coroutines for structure and streams; combine with asCoroutineDispatcher() |
Story Closing
The Kotlin feed went to staging the next morning, with one catch that logged the payload and greyed out the tile. When the Mumbai gateway sent another 12O that afternoon, the tile went grey, the log line quoted the bad reading, and the gateway team had a ticket before anyone in the ops room noticed.
The board went into a shared library, and the first team to use it was the Java one. Their build compiled; their code didn’t read well. The feed was StockBoardKt.lowStockFeed(board), a static method on a class named after a file. Constants lived behind Companion accessors. And in the pricing module they also depended on, a method called shelfLabel-C-dwYyA couldn’t be called from Java at all.
Their lead sent a screenshot and one line: “Is this what Kotlin libraries look like from the outside?”
“Only the ones nobody wrote for Java,” Lena said.
In Part 12, Kabir learns to make Kotlin code bilingual: what Java sees, the annotations that fix it, the Gradle plugins Spring needs, and the testing and quality tools that keep a Kotlin codebase honest.
This is Part 11 of a 16-part series: “Kotlin for Java Survivors: Life After Semicolons.”