skip to content

Why might collectLatest fail to cancel an in-flight block, and how do you make CPU-bound work in the block actually cancellable?

level: middleimportance: must knowfreq 55%

answer

  1. Cancellation is cooperative — needs a checkpoint
  2. Checkpoints: delay/yield/withContext/suspend + ensureActive()/isActive
  3. CPU loops are uncancellable without ensureActive()/yield()
  4. isActive only reads; ensureActive() throws
  5. Never swallow CancellationException

basics

~20 s

Cancellation 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 s

Coroutine 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 lines
kotlin
import 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

for a junior

Recognizes a pause point is needed for cancellation to work.

for a middle

Names ensureActive()/yield()/isActive and rewrites a CPU loop to cooperate.

for a senior

Distinguishes read-only isActive from throwing ensureActive(), and protects CancellationException from broad catches.

for a principal

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

context