This article explains how to build a retry with a guaranteed attempt count using Kotlin Flow's
retryWhen. The pattern wraps a one-shot suspend call asflow { } → retryWhen { } → single()and splits success from failure withResult.fold. It gives developers who are new to Flow the background they need to read or write this pattern, and diagrams the behavior for each scenario.
Baseline: Kotlin 2.0.0, kotlinx.coroutines 1.8.1 (checked 2026-09). The library code quoted here is copied from the 1.8.1 source. The example was compiled against this version and verified by running 10 scenarios.
Say you have a suspend function that calls an external HTTP API. On a transient failure (503, dropped connection), you want to try again. There are three requirements.
Retry-After header saysYou could put several exit conditions inside a while (true) loop. But deleting just one exit condition by mistake makes the loop never end. The pattern below gathers the termination condition into one place: a comparison against the attempt count.
Take a service that calls a weather forecast API. We model the real HTTP client as an interface.
data class Forecast(val city: String, val summary: String)
// Exception the client throws on a 4xx or 5xx response
class HttpStatusException(val status: Int, val retryAfterSeconds: Long?) : RuntimeException("HTTP $status")
interface ForecastApi {
// Throws IOException if no response is received at all
suspend fun fetch(requestId: String, city: String): Forecast
}
data class RetryPolicy(val maxAttempts: Int) {
companion object {
val NO_RETRY = RetryPolicy(maxAttempts = 1)
}
}
sealed interface FetchResult {
val attempts: Int
data class Success(override val attempts: Int, val forecast: Forecast) : FetchResult
data class Failure(override val attempts: Int, val status: Int) : FetchResult
}
The response status decides whether to retry and how long to wait.
private const val MAX_RETRY_AFTER_SECONDS = 25L
// Returns the seconds to wait if a retry is allowed, or null if not
fun retryDelaySeconds(e: HttpStatusException): Long? {
return when (e.status) {
503 -> (e.retryAfterSeconds ?: 0L).coerceAtMost(MAX_RETRY_AFTER_SECONDS)
502 -> 0L
else -> null
}
}
Here is the core function.
class ForecastService(private val api: ForecastApi) {
suspend fun fetchForecast(city: String, policy: RetryPolicy): FetchResult {
var attempts = 0
return runCatchingCancellable {
flow {
attempts++
emit(api.fetch(requestId = "forecast:$city:$attempts", city = city))
}.retryWhen { cause, retriedCount ->
if (retriedCount + 1 >= policy.maxAttempts) {
return@retryWhen false
}
val delaySeconds = when (cause) {
is IOException -> 0L
is HttpStatusException -> retryDelaySeconds(cause)
else -> null
} ?: return@retryWhen false
delay(delaySeconds.seconds)
true
}.single()
}.fold(
onSuccess = { forecast -> FetchResult.Success(attempts = attempts, forecast = forecast) },
onFailure = { e ->
if (e !is HttpStatusException) {
throw e
}
FetchResult.Failure(attempts = attempts, status = e.status)
},
)
}
}
runCatchingCancellable is not a standard library function. It is a helper that teams often define themselves to avoid the problem where the standard runCatching also swallows CancellationException (kotlinx.coroutines#1814).
inline fun <R> runCatchingCancellable(block: () -> R): Result<R> {
return try {
Result.success(block())
} catch (e: CancellationException) {
throw e
} catch (e: Throwable) {
Result.failure(e)
}
}
There are two key points.
retryWhen is what retries. When the predicate returns true, retryWhen collects the upstream flow again from the start. So the flow { } block runs one more time.retriedCount + 1 >= maxAttempts, the predicate always returns false. So the number of calls never exceeds maxAttempts.The flow { } builder does not run its block right away. It runs the block when someone collects (collect), and collecting again runs the block again from the start. This kind of flow is called a cold flow.
val numbers = flow {
println("block start")
emit(1)
}
numbers.collect { println(it) } // "block start", 1
numbers.collect { println(it) } // "block start", 1 ← the block runs again
This property is the basis of the retry. retryWhen does not rewind the failed call. It only collects the upstream flow anew. A new collect runs the block from the start, and the api.fetch inside it is called again. This is also why attempts in the example goes up. retryWhen does not increment the count. The block runs again, and attempts++ runs with it.
Operators fall into intermediate and terminal operators. An intermediate operator such as retryWhen only builds and returns a new flow and runs nothing. You must call a terminal operator such as single() or collect() to set the whole chain in motion. In the example, the .single() at the end of the chain is what actually starts the call.
The kotlinx.coroutines 1.8.1 implementation is short. Here it is verbatim.
public fun <T> Flow<T>.retryWhen(predicate: suspend FlowCollector<T>.(cause: Throwable, attempt: Long) -> Boolean): Flow<T> =
flow {
var attempt = 0L
var shallRetry: Boolean
do {
shallRetry = false
val cause = catchImpl(this)
if (cause != null) {
if (predicate(cause, attempt)) {
shallRetry = true
attempt++
} else {
throw cause
}
}
} while (shallRetry)
}
The behavior comes down to five points.
predicate(cause, attempt). On true, it adds 1 to attempt and collects againpredicate returns false, it throws the original exception cause as is, without wrapping itattempt counts from 0 (it is 0 on the first failure)predicate is a suspend function, so waiting with delay inside it and then returning true gives "wait, then retry"Two kinds of exceptions are not caught by retryWhen. The internal catchImpl rethrows both.
single() or by the collecting codeCancellationException raised when a coroutine is cancelledSo when the calling coroutine is cancelled, the cancellation goes through right away instead of being caught by the retry. The example's runCatchingCancellable also rethrows the cancellation exception, so a cancellation is never turned into a Failure result and swallowed.
retry(retries) is a thin wrapper around retryWhen. In 1.8.1 it is one line: retryWhen { cause, attempt -> attempt < retries && predicate(cause) }. The example uses retryWhen directly because the wait time differs by exception type and response.
The
attemptinretryWhenand theretriesinretry(retries)are the retry count. The total attempt count is that value plus 1. This is why the example comparesretriedCount + 1withmaxAttempts.
single() is a terminal operator. It collects the flow and returns one value. According to the KDoc, it throws NoSuchElementException for an empty flow and IllegalArgumentException if there is more than one element.
The example's flow { } block calls emit at most once per run. A failed attempt throws before emit, so it emits nothing. So the value single() receives is the one response from the last, successful attempt. If upstream ends with an exception in the end, single() throws that exception as is.
Result.fold(onSuccess, onFailure) calls onSuccess on success and onFailure on failure, and returns the value they produce. Both lambdas are inline, so a throw inside onFailure sends the exception straight out of the function.
The example uses this property to sort exceptions into two groups.
| Exception that failed last | Meaning | Result |
|---|---|---|
HttpStatusException |
The server responded with 4xx or 5xx | Returned as Failure(status) |
IOException |
No response was received | Exception propagates |
| Any other exception | Unexpected failure | Exception propagates |
CancellationException |
Coroutine cancelled | runCatchingCancellable rethrows it, so the cancellation proceeds as is |
When a response arrived but it was a failure, the caller gets it as a result and branches on it. When there was no response at all, the failure is signaled as an exception.
The retriedCount passed to the predicate counts from 0. attempts goes up by 1 each time the flow { } block runs. When the predicate is called, attempts = retriedCount + 1 always holds. Here is a trace with maxAttempts = 2.
| Moment | attempts |
requestId suffix | retriedCount |
retriedCount + 1 >= 2 |
Result |
|---|---|---|---|---|---|
| 1st call fails | 1 | :1 |
0 | false | Collects again if the failure is retryable |
| 2nd call fails | 2 | :2 |
1 | true | false: throws the exception as is |
A 3rd call never happens. The end is decided by the attempt-count comparison, not by an exit condition inside a loop.
The two requests have different requestIds, :1 and :2. You can use this suffix in the server logs to tell apart the retries of one call.
A failure that came with a response returns as a Failure result, not an exception, even when it fails to the end. If IOException occurs twice, it takes the same path, then splits at fold and the exception propagates.
This is a test written with JUnit 5 and AssertJ. runTest uses virtual time, so it does not actually wait for delay(7 seconds). Measuring the elapsed virtual time with testTimeSource makes the test fail if the wait logic is missing.
class ForecastServiceTest {
// Fake API that runs the prepared behavior in order, one per call
private class ScriptedApi(private val steps: List<() -> Forecast>) : ForecastApi {
private var calls = 0
override suspend fun fetch(requestId: String, city: String): Forecast {
return steps[calls++]()
}
}
@OptIn(ExperimentalCoroutinesApi::class)
@Test
fun `on 503 it waits for Retry-After, retries, and succeeds`() = runTest {
// given: 1st call 503 (Retry-After 7), 2nd call succeeds
val forecast = Forecast("Seoul", "Sunny")
val api = ScriptedApi(listOf({ throw HttpStatusException(503, retryAfterSeconds = 7) }, { forecast }))
val startMark = testTimeSource.markNow()
// when
val result = ForecastService(api).fetchForecast("Seoul", RetryPolicy(maxAttempts = 2))
// then
assertThat(startMark.elapsedNow()).isEqualTo(7.seconds)
assertThat(result).isEqualTo(FetchResult.Success(attempts = 2, forecast = forecast))
}
}
testTimeSourceandtestScheduler.currentTimeare@ExperimentalCoroutinesApiin 1.8.1. Without the opt-in, the compiler emits a warning. Add it only to tests that read virtual time.
attempts goes up when the flow block runs again. retryWhen does not raise itfalse throws the original exception as is. It is not wrapped, so fold can split by type with an is checkdelay inside the predicate. This works because the predicate is suspendretryWhen does not catch it, and runCatchingCancellable rethrows itemit a value twice, single() throws IllegalArgumentException. When you edit the flow block, check that it calls emit only onceattempts is used within a single call only. The chain runs sequentially, so there is no concurrent access. Do not change it so that several calls share the same variableretry(n) { cause -> ... }: When the wait time is fixed and the decision looks only at the exception type, this is shorter. You can also use delay inside the predicate lambdafor (attempt in 1..maxAttempts) loop: More familiar to teammates who do not know Flow. But you need a separate variable outside the loop to hold the "last failure" to returnflow + retryWhen + single): The termination condition, the retry decision, and the wait all sit in one predicate