本文介绍如何用 Kotlin Flow 的
retryWhen实现尝试次数有保证的重试。这个模式用flow { } → retryWhen { } → single()包装一次性的 suspend 调用,再用Result.fold区分成功与失败。文中整理了初次接触 Flow 的开发者阅读或自己编写这个模式时需要的背景知识,并用图示说明各场景下的行为。
基准版本:Kotlin 2.0.0,kotlinx.coroutines 1.8.1(2026-09 确认)。正文引用的库代码摘自 1.8.1 源码,示例已在该版本上编译,并运行了 10 个场景加以验证。
有一个调用外部 HTTP API 的 suspend 函数。遇到暂时性失败(503、连接中断)时,希望重新尝试。条件有三个。
Retry-After 头指定的时间也可以在 while (true) 里放多个退出条件。但只要误删其中一个退出条件,就会变成永不结束的循环。下面的模式把结束条件集中到一处,即比较尝试次数。
假设有一个调用天气预报 API 的服务。这里用接口代替真实的 HTTP 客户端。
data class Forecast(val city: String, val summary: String)
// 收到 4xx、5xx 响应时客户端抛出的异常
class HttpStatusException(val status: Int, val retryAfterSeconds: Long?) : RuntimeException("HTTP $status")
interface ForecastApi {
// 完全收不到响应时抛出 IOException
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
}
是否重试以及等待时间,由响应状态决定。
private const val MAX_RETRY_AFTER_SECONDS = 25L
// 可以重试则返回要等待的秒数,否则返回 null
fun retryDelaySeconds(e: HttpStatusException): Long? {
return when (e.status) {
503 -> (e.retryAfterSeconds ?: 0L).coerceAtMost(MAX_RETRY_AFTER_SECONDS)
502 -> 0L
else -> null
}
}
这是核心函数。
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 不是标准库函数。标准的 runCatching 会连 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)
}
}
关键有两点。
retryWhen。 判断函数返回 true 时,retryWhen 会从头重新收集上游 flow。因此 flow { } 代码块会再执行一次。retriedCount + 1 >= maxAttempts 时,判断函数一定返回 false。因此调用次数不会超过 maxAttempts。flow { } 构建器不会立即执行代码块。有人收集(collect)时才执行,再次收集时会从头重新执行代码块。这样的 flow 称为 cold flow。
val numbers = flow {
println("代码块开始")
emit(1)
}
numbers.collect { println(it) } // "代码块开始", 1
numbers.collect { println(it) } // "代码块开始", 1 ← 代码块再次执行
这个特性是重试的基础。retryWhen 不会回退已经失败的调用,只是重新收集上游 flow。重新收集时,代码块从头再跑,其中的 api.fetch 会被再次调用。示例中 attempts 增加的原因也一样。不是 retryWhen 在增加计数,而是代码块重新执行时,attempts++ 又跑了一遍。
操作符分为中间操作符和终端操作符。retryWhen 这类中间操作符只会生成并返回新的 flow,不执行任何东西。必须调用 single()、collect() 这类终端操作符,整条链才会运转。示例中真正启动调用的是链末尾的 .single()。
kotlinx.coroutines 1.8.1 的实现很短。这里原样摘录。
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)
}
其行为可归纳为五点。
predicate(cause, attempt);返回 true 就把 attempt 加 1 并重新收集predicate 返回 false 时,原来的异常 cause 不会被包装,而是原样抛出attempt 从 0 开始计数(第一次失败时为 0)predicate 是 suspend 函数,所以可以在里面用 delay 等待后再返回 true,就实现了“等待后重试”有两类异常是 retryWhen 不会捕获的。内部的 catchImpl 会把它们原样抛出。
single() 或收集端代码抛出的异常CancellationException因此调用方协程被取消时,不会被重试拦住,而是立即取消。示例中的 runCatchingCancellable 也会重新抛出取消异常,所以取消不会被转换成 Failure 结果而被吞掉。
retry(retries) 是对 retryWhen 的一层薄封装。1.8.1 的实现只有一行:retryWhen { cause, attempt -> attempt < retries && predicate(cause) }。示例中异常类型不同、响应不同,等待时间也不同,所以直接使用 retryWhen。
retryWhen的attempt和retry(retries)的retries是重试次数。总尝试次数是它加 1。这就是示例把retriedCount + 1与maxAttempts比较的原因。
single() 是终端操作符。它收集 flow 并返回一个值。根据 KDoc,flow 为空时抛出 NoSuchElementException,元素多于一个时抛出 IllegalArgumentException。
示例的 flow { } 代码块每执行一次,最多 emit 一次。失败的尝试在 emit 之前就出了异常,不会发出值。所以 single() 收到的是最后一次成功尝试的那一个响应。如果上游最终以异常结束,single() 也会原样抛出该异常。
Result.fold(onSuccess, onFailure) 在成功时调用 onSuccess,失败时调用 onFailure,并返回其返回值。两个 lambda 都是 inline,所以在 onFailure 里 throw 时,异常会直接从函数外抛出。
示例利用这个特性把异常分成两类。
| 最终失败的异常 | 含义 | 结果 |
|---|---|---|
HttpStatusException |
服务器以 4xx、5xx 响应 | 以 Failure(status) 返回 |
IOException |
没有收到响应 | 异常向外传播 |
| 其他异常 | 意料之外的失败 | 异常向外传播 |
CancellationException |
协程取消 | runCatchingCancellable 重新抛出,取消照常进行 |
收到了响应但失败的情况,调用方作为结果接收并分支处理;完全没有响应的情况,则以异常通知。
传给判断函数的 retriedCount 从 0 开始计数。attempts 则在每次 flow { } 代码块执行时加 1。判断函数被调用的时刻,始终满足 attempts = retriedCount + 1。以 maxAttempts = 2 为例来跟踪,如下所示。
| 时刻 | attempts |
requestId 结尾 | retriedCount |
retriedCount + 1 >= 2 |
结果 |
|---|---|---|---|---|---|
| 第 1 次调用失败 | 1 | :1 |
0 | 假 | 若为重试对象则重新收集 |
| 第 2 次调用失败 | 2 | :2 |
1 | 真 | false: 原样抛出异常 |
不会发生第 3 次调用。结束不是由循环里的退出条件决定,而是由尝试次数的比较决定。
两次请求的 requestId 分别以 :1、:2 结尾,互不相同。在服务器日志里可以用这个后缀区分同一次调用的各次重试。
有响应的失败,即使一直失败到最后,也不是以异常而是以 Failure 结果返回。IOException 出现两次时走同一条路径,到 fold 处分流,异常向外传播。
这是用 JUnit 5 和 AssertJ 编写的测试。runTest 使用虚拟时间,不会真的等待 delay(7秒)。用 testTimeSource 测量流逝的虚拟时间,等待逻辑一旦缺失,测试就会失败。
class ForecastServiceTest {
// 每次调用依次执行预先准备好的行为的假 API
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 `503时按 Retry-After 等待后重试并成功`() = runTest {
// given: 第 1 次 503(Retry-After 7), 第 2 次成功
val forecast = Forecast("Seoul", "晴")
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))
}
}
testTimeSource和testScheduler.currentTime在 1.8.1 中是@ExperimentalCoroutinesApi。没有 opt-in 会出现编译警告。只在读取虚拟时间的测试上加这个注解。
attempts 在 flow 代码块重新执行时增加。 不是 retryWhen 在增加这个值false 时,原来的异常会原样抛出。 因为没有被包装,所以可以在 fold 中用 is 检查区分类型delay。 因为判断函数是 suspend,所以才能这样做retryWhen 不会捕获它,runCatchingCancellable 也会重新抛出emit 两次值时,single() 会抛出 IllegalArgumentException。 修改 flow 代码块时,要确认 emit 只执行一次attempts 只在一次调用内部使用。 因为链是顺序执行的,所以不存在并发访问。不能改成让多次调用共享同一个变量retry(n) { cause -> ... }:如果等待时间固定,判断只看异常类型,这个写法更短。判断 lambda 里也可以使用 delayfor (attempt in 1..maxAttempts) 循环:对不了解 Flow 的团队成员更熟悉。不过需要在循环之外另外持有一个变量,用来返回“最后一次失败”flow + retryWhen + single):结束条件、重试判断、等待集中在判断函数这一处