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 / ReactorKotlinNote
Flux<T> (cold)Flow<T>Built on suspend; no Publisher protocol
Mono<T>suspend fun …(): TOne value doesn’t need a stream type
Flux.generate(…) / a computed sequenceflow { emit(…) }Sequential code with emit
Flux.create(sink -> listener …)callbackFlow { … awaitClose { } }Bridges a Java listener API
flatMap(x -> asyncCall(x))map { suspendCall(it) }, or flatMapMergemap is sequential and ordered; flatMapMerge is concurrent (16 by default)
concatMapflatMapConcatOpt-in (ExperimentalCoroutinesApi), like flatMapMerge
onBackpressureBuffer / onBackpressureLatestbuffer() / conflate()Back-pressure is suspension
subscribeOnflowOn(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()MutableSharedFlowBroadcast; replay and buffer configurable
BlockingQueue between threadsChannelEach element to one receiver; suspends instead of blocking
StepVerifier + virtual timerunTest + 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

TypeWho producesWho receivesValue when nobody listensJava/Reactor analogue
Flow (cold)The collector, by collectingThat collector onlyNothing runsCold Flux
SharedFlow (hot)Code that calls emitEvery current subscriberDropped, unless replaySinks.many().multicast().directBestEffort()
StateFlow (hot)Code that sets valueEvery subscriber, latest value firstKept: there is always a current valueSinks.many().replay().latestOrDefault(x)
ChannelsendExactly one receiver per elementBuffered or suspendedBlockingQueue

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.Dispatchers
import kotlinx.coroutines.delay
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.filter
import kotlinx.coroutines.flow.flow
import kotlinx.coroutines.flow.flowOn
import kotlinx.coroutines.flow.map
import kotlinx.coroutines.flow.take
import kotlinx.coroutines.flow.toList
import kotlinx.coroutines.runBlocking
import 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.Flow
import kotlinx.coroutines.flow.buffer
import kotlinx.coroutines.flow.flow
import 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:

graph LR subgraph P["producer coroutine (flowOn's dispatcher, if any)"] A["flow { emit }"] --> B["mapped"] end B -->|"send: suspends when full"| C[("channel: 64 slots by default")] C -->|receive| D subgraph K["collector coroutine (the caller's context)"] D["collect { }"] end

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.ExperimentalCoroutinesApi
import kotlinx.coroutines.delay
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.flatMapConcat
import kotlinx.coroutines.flow.flatMapMerge
import kotlinx.coroutines.flow.flow
import kotlinx.coroutines.flow.flowOf
import kotlinx.coroutines.flow.map
import kotlinx.coroutines.flow.toList
import 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.0
fun 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.ExperimentalCoroutinesApi
import kotlinx.coroutines.FlowPreview
import kotlinx.coroutines.delay
import kotlinx.coroutines.flow.buffer
import kotlinx.coroutines.flow.collectLatest
import kotlinx.coroutines.flow.conflate
import kotlinx.coroutines.flow.debounce
import kotlinx.coroutines.flow.flow
import kotlinx.coroutines.flow.toList
import kotlinx.coroutines.test.TestScope
import kotlinx.coroutines.test.currentTime
import 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.0
suspend 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.0
fun 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.delay
import kotlinx.coroutines.flow.MutableSharedFlow
import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.flow.SharingStarted
import kotlinx.coroutines.flow.asStateFlow
import kotlinx.coroutines.flow.flow
import kotlinx.coroutines.flow.stateIn
import kotlinx.coroutines.flow.take
import kotlinx.coroutines.flow.toList
import kotlinx.coroutines.launch
import 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.ExperimentalCoroutinesApi
import kotlinx.coroutines.channels.Channel
import kotlinx.coroutines.channels.produce
import kotlinx.coroutines.coroutineScope
import kotlinx.coroutines.delay
import kotlinx.coroutines.launch
import kotlinx.coroutines.runBlocking
import 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.0
fun 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.async
import kotlinx.coroutines.channels.awaitClose
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.callbackFlow
import kotlinx.coroutines.flow.take
import kotlinx.coroutines.flow.toList
import kotlinx.coroutines.runBlocking
import 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.delay
import kotlinx.coroutines.flow.flow
import kotlinx.coroutines.flow.map
import kotlinx.coroutines.flow.toList
import kotlinx.coroutines.reactive.asFlow
import kotlinx.coroutines.reactive.awaitSingle
import kotlinx.coroutines.reactor.awaitSingleOrNull
import kotlinx.coroutines.reactor.asFlux
import kotlinx.coroutines.reactor.flux
import kotlinx.coroutines.reactor.mono
import kotlinx.coroutines.runBlocking
import reactor.core.publisher.Flux
import 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.asCoroutineDispatcher
import kotlinx.coroutines.async
import kotlinx.coroutines.awaitAll
import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.withContext
import 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.

SituationChooseWhy
Existing Spring MVC + JPA/JDBC service, blocking libraries throughoutVirtual 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 WebFluxCoroutines (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-nothingCoroutinesStructured concurrency is the point; Java’s StructuredTaskScope is still a preview API
Streams of values, live updates, back-pressureCoroutines (Flow)Virtual threads have no stream model
A blocking SDK inside coroutine codeCoroutines + a virtual-thread (or IO) dispatcherStructure from coroutines, cheap blocking from the JVM
CPU-bound workNeither helps; Dispatchers.Default or a bounded poolWaiting isn’t the bottleneck
Java-heavy team, Kotlin only at the edgesVirtual threadsNo 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.0
fun 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.ShelfScan
import com.shelfwise.part11.dashboard.StockBoard
import kotlinx.coroutines.ExperimentalCoroutinesApi
import kotlinx.coroutines.delay
import kotlinx.coroutines.flow.first
import kotlinx.coroutines.launch
import kotlinx.coroutines.test.advanceTimeBy
import kotlinx.coroutines.test.runCurrent
import kotlinx.coroutines.test.runTest
import org.junit.jupiter.api.Assertions.assertEquals
import org.junit.jupiter.api.Test
@OptIn(ExperimentalCoroutinesApi::class) // advanceTimeBy, runCurrent, currentTime
class 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 stateIn and shareIn. A field private val scope = CoroutineScope(SupervisorJob() + Dispatchers.Default) and a @PreDestroy fun close() = scope.cancel() tie the shared upstream to the bean’s lifecycle. GlobalScope would outlive the application context, and a request’s scope would stop the upstream when the first request ends.

Tip — runTest belongs in tests. This part’s examples call it from main only to make timings deterministic. Production code collects in a real scope: a controller method, a bean scope, or runBlocking at the edge of a command-line tool.

Gotcha — porting flatMap changes either the order or the latency. A Reactor .flatMap { fetch(it) } runs up to 256 inner calls at once. Rewritten as map { fetch(it) }, the same chain makes one call at a time: correct, ordered, and as slow as the sum of the calls. flatMapMerge restores concurrency, but its default limit is 16 (the system property kotlinx.coroutines.flow.defaultConcurrency changes it), and results arrive in completion order. Decide which you need per chain, and pass concurrency explicitly when it matters. Reactor’s flatMapSequential (concurrent, but in source order) has no built-in equivalent.

Gotcha — mutating the object inside a StateFlow emits 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 a StateFlow and change them with update { }.

Gotcha — tryEmit is not Reactor’s tryEmitNext. Reactor’s Sinks.many().multicast().directBestEffort() reports FAIL_ZERO_SUBSCRIBER when nobody is listening. A MutableSharedFlow without a buffer does the opposite: tryEmit returns true when nobody is subscribed, because dropping the value succeeded, and false as soon as someone is, because delivering it would mean suspending (both lines in HotFlows.kt). Code ported from Reactor that treats true as “delivered” is wrong in both directions. Use the suspending emit, or give the flow a buffer (extraBufferCapacity) and handle false.

Gotcha — a try/catch around emit catches the collector’s exceptions. emit runs the collector’s code. A Java-style try { 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 the catch { } operator, which only sees upstream failures:

import kotlinx.coroutines.flow.catch
import kotlinx.coroutines.flow.flow
import 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.Flow is 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 want kotlinx.coroutines.flow.Flow; if you really have a JDK Publisher, kotlinx-coroutines-jdk9 converts it with asFlow().

Gotcha — a Flow re-runs for every collector. A Java Stream can be consumed once; a cold Flow can be collected any number of times, and each collection runs the producer again (“scanner connected” twice above). Two dashboards collecting stockLevels("PUN-014") open two scanner connections. Share it with stateIn or shareIn.


Debugging and Troubleshooting

SymptomLikely causeFix
A stream stops updating, with no error anywhereAn error handler that completes the stream: Reactor’s onErrorResume { Flux.empty() }, or a catch { } that emits nothingLog in the handler and emit a visible fallback, or let it fail and restart deliberately with retry/retryWhen
IllegalStateException: Flow invariant is violatedemit called from a different coroutine context inside flow { }, usually via withContextUse flowOn for the upstream, or channelFlow { } if you really need several producers
A collector never returnsCollecting a StateFlow or SharedFlow: they never completefirst { }, take(n), or collect in a scope you cancel
A test hangs, then fails with a timeoutAn infinite collector launched in the runTest scopeLaunch it in backgroundScope
Dashboard misses intermediate valuesStateFlow and conflate() keep only the latest, by designSharedFlow with a buffer, if every value matters
Virtual time doesn’t move in a testThe code under test uses its own Dispatchers.Default/IOInject the dispatcher or scope; pass the test’s
NoSuchElementException from awaitSingle()The Mono was emptyawaitSingleOrNull()
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

ConceptRemember
FlowCold, sequential, suspend-based; each collection runs the producer
OperatorsA flow that collects its upstream; emit calls the collector directly
Back-pressureSuspension at emit; buffer, conflate, collectLatest, debounce (preview)
flatMap habitsmap { suspendCall() } serialises; flatMapMerge reorders, 16 at a time by default; choose per chain
buffer / flowOnSplit the pipeline into two coroutines and a channel; flowOn moves only what’s above it; never withContext around emit
StateFlowObservable value; skips equal values; immutable data changed with update { }
SharedFlowHot event broadcast; replay decides what late subscribers see
stateIn / shareInShare one upstream between many collectors, in a scope you own
ChannelA suspending queue: each element to one receiver; for distributing work
callbackFlowJava listener → Flow: trySend in the callback, awaitClose { unregister } at the end
Reactor interopawaitSingle, asFlow, mono { }, flux { }, asFlux
TestingrunTest virtual time; infinite collectors in backgroundScope; inject dispatchers
Virtual threads vs coroutinesVirtual 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.”