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、接続切断)なら再試行したい。条件は3つある。
Retry-Afterヘッダーの分だけ待つwhile (true)の中に脱出条件を複数置く方法もある。しかし脱出条件を1つ誤って消しただけで、終わらないループになる。以下のパターンは終了条件を試行回数の比較1か所にまとめる。
天気予報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)
}
}
要点は2つある。
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)
}
動作は5つにまとめられる。
predicate(cause, attempt)を呼び、trueならattemptを1増やしてから再収集するpredicateがfalseなら、元の例外causeを包まずそのまま投げるattemptは0から数える(最初の失敗のときは0)predicateはsuspend関数なので、中でdelayで待ってからtrueを返せば「待機後に再試行」になるretryWhenが捕まえない例外が2種類ある。内部のcatchImplがこの2つをそのまま投げる。
single()や収集側のコードが投げた例外CancellationExceptionそのため、呼び出したコルーチンがキャンセルされると、再試行で引き止めずにすぐキャンセルされる。例のrunCatchingCancellableもキャンセル例外を再びスローするので、キャンセルがFailure結果に変わって飲み込まれることはない。
retry(retries)はretryWhenを薄く包んだものである。1.8.1の実装はretryWhen { cause, attempt -> attempt < retries && predicate(cause) }の1行である。例は例外の種類と応答ごとに待ち時間が違うので、retryWhenを直接使う。
retryWhenのattemptとretry(retries)のretriesは再試行回数である。全体の試行回数はこれに1を足した値になる。例がretriedCount + 1とmaxAttemptsを比較する理由である。
single()は終端演算子である。flowを収集して値を1つ返す。KDocによると、空のflowならNoSuchElementExceptionを、要素が2つ以上ならIllegalArgumentExceptionを投げる。
例のflow { }ブロックは、1回実行されるごとにemitを最大1回行う。失敗した試行はemitの前に例外が起きるので、値を出さない。そのためsingle()が受け取る値は成功した最後の試行の応答1つである。上流が最後まで例外で終わったときは、single()もその例外をそのまま投げる。
Result.fold(onSuccess, onFailure)は、成功ならonSuccessを、失敗ならonFailureを呼んで、その戻り値を返す。どちらのラムダもinlineなので、onFailureの中でthrowすると、関数の外へ例外がそのまま出る。
例はこの性質を使って、例外を2種類に分ける。
| 最後まで失敗した例外 | 意味 | 結果 |
|---|---|---|
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回目の呼び出しは起きない。終了はループ内の脱出条件ではなく、試行回数の比較で決まる。
2つのリクエストのrequestIdは:1、:2で異なる。サーバーログでは、この接尾辞で同じ呼び出しの再試行を区別できる。
応答のある失敗は、最後まで失敗しても例外ではなくFailure結果として返る。IOExceptionが2回起きた場合は同じ経路を進み、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が1回だけか確認する必要があるattemptsは1回の呼び出しの中だけで使う。 チェーンが逐次で動くので、同時アクセスはない。複数の呼び出しが同じ変数を共有するように変えてはならないretry(n) { cause -> ... }: 待ち時間が固定で、判定が例外の種類だけを見るなら、こちらのほうが短い。判定ラムダの中でdelayも使えるfor (attempt in 1..maxAttempts)ループ: Flowを知らないチームメンバーにはこちらのほうが馴染みがある。ただしループの外で「最後の失敗」を返すための変数を別に持つ必要があるflow + retryWhen + single): 終了条件・再試行の判定・待機が、判定関数の1か所にまとまる