Why might collectLatest fail to cancel an in-flight block, and how do you make CPU-bound work in the block actually cancellable?
answer
- Cancellation is cooperative — needs a checkpoint
- Checkpoints: delay/yield/withContext/suspend + ensureActive()/isActive
- CPU loops are uncancellable without ensureActive()/yield()
- isActive only reads; ensureActive() throws
- Never swallow CancellationException
basics
~20 sCancellation only happens when the running code reaches a pause point. If your block does heavy non-stop computation, there's no pause point, so the old block keeps running. Add checks like ensureActive() or yield() so it can stop.
solid answer
~40 sCoroutine cancellation is cooperative: a coroutine reacts to cancellation only at a suspension point (delay, yield, withContext, any suspend call) or when it inspects its CancellationException state via isActive / ensureActive(). collectLatest cancels the previous action coroutine, but a tight CPU loop with no suspension and no cooperation will ignore that cancel signal and run to completion, defeating the 'only latest completes' guarantee. To fix CPU-bound blocks: periodically call ensureActive() (throws CancellationException if cancelled) or yield() (also a suspension point) inside the loop, or move the heavy work onto a dispatcher and check isActive. Note isActive only *reads* the flag — you still must break out yourself. ensureActive()/yield() are the idiomatic, throw-on-cancel choices.
code
kotlin · 16 linesimport kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
suspend fun cancellableWork(n: Int): Long = coroutineScope {
var acc = 0L
for (i in 0 until 50_000_000) {
if (i % 100_000 == 0) ensureActive() // cooperate
acc += (i.toLong() * n)
}
acc
}
fun main() = runBlocking {
flowOf(1, 2, 3).collectLatest { n -> println(cancellableWork(n)) }
// only the result for 3 prints, because work for 1 and 2 is cancelled
}go deeper
Recognizes a pause point is needed for cancellation to work.
Names ensureActive()/yield()/isActive and rewrites a CPU loop to cooperate.
Distinguishes read-only isActive from throwing ensureActive(), and protects CancellationException from broad catches.
Sets a team convention for cancellation checkpoints in hot loops and reasons about the granularity/overhead tradeoff of the check frequency.
## Cooperative cancellation, defined Kotlin coroutines are **never force-killed**. Cancelling a coroutine sets a flag and, at the next *cancellation checkpoint*, a `CancellationException` is thrown to unwind it. A checkpoint is: - any **suspension point** from `kotlinx.coroutines` — `delay`, `yield`, `withContext`, `await`, or any `suspend` function that itself checks cancellation; **and** - explicit checks: `ensureActive()` (throws `CancellationException` if cancelled) or reading `isActive` (a `Boolean` you must act on yourself). A plain blocking/CPU loop has none of these, so it is **uncancellable**. ## How this bites collectLatest `collectLatest` implements 'cancel previous, start latest' by cancelling the child coroutine running your block. If that block is CPU-bound with no checkpoint, the cancel signal is ignored and the **old block finishes anyway** — so you may see stale results complete, and you lose the freshness guarantee. ```kotlin // BAD: old block can't be cancelled flow.collectLatest { data -> var acc = 0L for (i in 0 until 100_000_000) acc += heavyPure(i, data) // no checkpoint! render(acc) } ``` ## Making it cancellable ```kotlin // GOOD: cooperate with cancellation flow.collectLatest { data -> var acc = 0L for (i in 0 until 100_000_000) { if (i % 10_000 == 0) ensureActive() // throws if superseded acc += heavyPure(i, data) } render(acc) } ``` Alternatives: - **`yield()`** inside the loop — both a suspension point and a cancellation check, but it also reschedules, adding overhead. - **`ensureActive()`** — cheapest explicit check; throws immediately if cancelled; doesn't yield the thread. - **`isActive`** — `while (isActive) { ... }`: you read the flag and `break`/return yourself; no exception thrown. - Offload to `Dispatchers.Default` via `withContext` and check periodically; `withContext` itself is a suspension point at entry/exit but not during the loop body. ## Key nuance `CancellationException` is special: it must propagate. Don't swallow it in a broad `catch (e: Exception)` — rethrow, or use `catch (e: CancellationException) { throw e }`, otherwise you break structured-concurrency cancellation and `collectLatest` semantics.
- What is the difference between ensureActive() and isActive?ensureActive() throws CancellationException immediately if the coroutine is cancelled; isActive is a Boolean you read and must act on (e.g. break) yourself.
- Why must you not catch and swallow CancellationException?It is the signal that powers cooperative cancellation and structured concurrency; swallowing it makes the block uncancellable and breaks collectLatest's freshness guarantee.
A runner who only checks their phone at water stations: if the race is called off mid-stretch, they keep running until the next station.
saying these in an interview costs you the question
- Believing cancellation forcibly interrupts CPU work
- Catching Exception broadly and silently absorbing CancellationException
- Thinking isActive throws on cancellation
- Adding withContext but doing the whole loop without any in-loop checkpoint
- Claiming you never need to do anything for CPU-bound blocks