How does flatMapLatest implement cancellation, and what must your inner flow do to honour it correctly?
answer
- Built on transformLatest
- Cancellation is cooperative, needs suspension points
- ensureActive()/yield() in CPU loops
- Clean up with try/finally or onCompletion
- Never swallow CancellationException
basics
~20 sWhen a new value arrives, flatMapLatest cancels the coroutine running the previous inner flow and starts a new one. Your inner flow must be cancellation-cooperative — use suspending calls or check for cancellation — so it actually stops.
solid answer
~40 sflatMapLatest is built on transformLatest. For each upstream value it launches the inner flow in a child coroutine; when the next value arrives it cancels that child via structured concurrency before starting the new one. Cancellation throws CancellationException at the next suspension point. Therefore the inner flow only stops promptly if it is cooperative: it must suspend (delay, network/IO suspend calls, ensureActive/yield) so the exception can propagate. A tight non-suspending CPU loop ignores cancellation and keeps running. Resources opened inside the inner flow should be released in flow {} via try/finally or onCompletion, because cancellation unwinds through them. Note CancellationException is normal control flow and must not be swallowed by a broad catch (e.g. catch (e: Exception)); rethrow it or use currentCoroutineContext().ensureActive().
code
kotlin · 14 linesimport kotlinx.coroutines.flow.*
import kotlinx.coroutines.delay
val keystrokes = flowOf("k", "ko", "kot")
fun search(q: String) = flow {
delay(200L) // suspension point -> cancellable
emit("results for $q")
}
suspend fun typeahead() {
keystrokes
.flatMapLatest { search(it) } // each new query cancels the prior search
.collect(::println) // only 'results for kot' survives
}go deeper
Knows flatMapLatest cancels the previous inner flow when a new value arrives.
Explains it builds on transformLatest and is the right tool for search-as-you-type.
Articulates cooperative cancellation, suspension points, ensureActive/yield, and not swallowing CancellationException.
Designs cancellation-safe resource handling, reasons about debounce+flatMapLatest pipelines and edge cases like fast emissions starving inner flows.
## What flatMapLatest does `flatMapLatest { transform }` re-subscribes on every upstream emission and **cancels the previous inner flow**. It is implemented as: ```kotlin public fun <T, R> Flow<T>.flatMapLatest( transform: suspend (value: T) -> Flow<R> ): Flow<R> = transformLatest { emitAll(transform(it)) } ``` So the mechanism is `transformLatest`: the block for the previous value is **cancelled** when a new value arrives, then the block runs for the new value. ## How cancellation works Kotlin coroutine cancellation is **cooperative**: - Cancelling a coroutine sets its job to cancelling and throws `CancellationException` at the next **suspension point**. - Suspension points: `delay`, suspending IO, `emit`, `yield`, `ensureActive()` checks. - A pure CPU loop with no suspension never observes cancellation and keeps burning until it finishes. ```kotlin // BAD: ignores cancellation fun bad(n: Int) = flow { var x = 0L while (x < 1_000_000_000L) x++ // no suspension -> can't be cancelled emit(n) } // GOOD: cooperative fun good(n: Int) = flow { repeat(1000) { ensureActive() // or yield() // ...heavy chunk... } emit(n) } ``` ## Resource cleanup Because cancellation unwinds the inner flow's coroutine, release resources in `try/finally`, or use `.onCompletion { }`: ```kotlin fun stream(id: Int) = flow { val conn = open(id) try { emitAll(conn.events()) } finally { conn.close() } } ``` ## Do not swallow CancellationException `CancellationException` signals normal cancellation. A broad `catch (e: Exception)` inside the inner flow will eat it and break structured concurrency. Rethrow it, or call `currentCoroutineContext().ensureActive()` after catching to re-throw if cancelled. (`catch` flow operator already rethrows CancellationException for you.) ## Practical use Search-as-you-type: every keystroke cancels the in-flight request: ```kotlin queries.debounce(300L) .flatMapLatest { q -> repository.search(q) } .collect { render(it) } ``` ## Key terms - **Cooperative cancellation**: code must reach a suspension point for cancellation to take effect. - **Suspension point**: a place where a coroutine can pause/resume (and thus check cancellation). - **Structured concurrency**: child coroutines are tied to parents; cancelling the parent cancels children.
- Why might an inner flow doing heavy CPU work not get cancelled by flatMapLatest?Cancellation is cooperative; without a suspension point (delay/yield/ensureActive) the CancellationException is never thrown, so the loop keeps running.
- What is the difference between flatMapLatest and mapLatest?mapLatest transforms each value to a plain value (cancelling the previous transform); flatMapLatest transforms to a Flow and emits all its values (cancelling the previous inner flow).
Like swapping the record on a turntable: the new track only really starts once the needle on the old one lifts — and it only lifts at a groove (a suspension point).
saying these in an interview costs you the question
- Thinking cancellation is preemptive/instant regardless of code
- Swallowing CancellationException with catch (e: Exception)
- Forgetting resource cleanup on cancellation
- Confusing flatMapLatest with flatMapMerge(concurrency=1)
- Believing flatMapLatest waits for the old inner flow to complete