You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在Kotlin Flow中实现持久收集,修复Ktor WebSocket重试机制问题?

问题描述

使用Kotlin Flow收集和发送Ktor WebSocket消息,需求是在httpClient.webSocket抛出异常时实现重试机制,并且后续能通过调用retryWith函数手动触发流代码块重新运行。

实际场景:无网络连接时,httpClient.webSocket抛出UnresolvedAddressException,重试块会执行3次重试;当用户连接网络后,调用retryWith函数期望重新运行流代码块,但此时流已完成,exceptionFlow无法被收集,导致无法触发重试。

尝试的代码如下:

fun initConnection(scope: CoroutineScope) = flow {
        coroutineScope {
            val exceptionHandler = CoroutineExceptionHandler { _, exception ->
                // Catch the exception and emit it to the shared flow.
                println("*** emiting exception $exception")
                exceptionFlow.tryEmit(exception)
            }
            launch(exceptionHandler) {
                exceptionFlow.collect {
                    println("*** throwing $it")
                    throw it //throw here to trigger retry
                }
            }

            httpClient.webSocket(
                block = {
                    receiveAll {
                        emit(it)
                        handleIncoming(it)
                    }
                },
            )
        }
    }.retry(3) { cause ->
        println("*** retrying cause: $cause")
        connectionState = ConnectionState.RETRYING

        delay(500L)

        true
    }.catch {
        println("*** caught $it")
        connectionState = ConnectionState.DISCONNECTED
        emit(Incoming.Exception(it))
        close()
    }.onCompletion {
        println("*** onCompletion $it")
    }


suspend fun retryWith(throwable: Throwable) {
    exceptionFlow.emit(throwable)
}
修复方案

核心问题是原流在3次重试失败后进入catch块并调用close(),导致流彻底终止,后续无法再触发重试。需要调整流的结构,让它保持活跃状态,同时用信号流控制重试时机。

关键修改点

  • 用**共享流(MutableSharedFlow)**作为外部手动重试的触发信号,替代原有的异常抛出触发方式,避免流提前终止。
  • 重构流的执行逻辑,用循环包裹WebSocket连接代码,失败后先执行自动重试,耗尽后等待外部信号再重试。
  • 分离自动重试与手动重试的逻辑,确保两种场景都能正确触发连接尝试。

修改后的完整代码

首先定义状态和触发流:

enum class ConnectionState {
    CONNECTING, CONNECTED, RETRYING, DISCONNECTED, WAITING_FOR_RETRY
}

// 外部触发重试的信号流,replay=0确保只发送给活跃的收集者
private val retryTrigger = MutableSharedFlow<Unit>(replay = 0)

var connectionState: ConnectionState = ConnectionState.DISCONNECTED
    private set

重构initConnection函数:

fun initConnection(scope: CoroutineScope) = flow {
    while (true) {
        connectionState = ConnectionState.CONNECTING
        try {
            httpClient.webSocket {
                connectionState = ConnectionState.CONNECTED
                // 持续接收WebSocket消息
                receiveAll {
                    emit(it)
                    handleIncoming(it)
                }
            }
            // WebSocket正常关闭时退出循环,可根据需求调整是否需要自动重连
            break
        } catch (e: Exception) {
            connectionState = ConnectionState.RETRYING
            var autoRetryCount = 0
            var connectionSuccess = false

            // 执行3次自动重试
            while (autoRetryCount < 3 && !connectionSuccess) {
                delay(500L)
                try {
                    httpClient.webSocket {
                        connectionState = ConnectionState.CONNECTED
                        receiveAll {
                            emit(it)
                            handleIncoming(it)
                        }
                    }
                    connectionSuccess = true
                } catch (retryE: Exception) {
                    autoRetryCount++
                    if (autoRetryCount >= 3) {
                        // 自动重试耗尽,进入等待手动重试状态
                        connectionState = ConnectionState.WAITING_FOR_RETRY
                        // 挂起直到收到手动重试信号
                        retryTrigger.first()
                    }
                }
            }
        }
    }
}.catch { e ->
    connectionState = ConnectionState.DISCONNECTED
    emit(Incoming.Exception(e))
}.onCompletion {
    connectionState = ConnectionState.DISCONNECTED
}

修改retryWith函数,直接发送重试信号:

suspend fun retryWith() {
    retryTrigger.emit(Unit)
}

逻辑说明

  1. 外层while(true)让流保持活跃,除非WebSocket正常关闭或遇到无法处理的异常。
  2. 自动重试阶段:连接失败后,最多尝试3次自动重连,每次间隔500ms。
  3. 手动重试阶段:3次自动重试失败后,流进入挂起状态,等待retryTrigger的信号;调用retryWith时发送信号,流会立即重新尝试建立连接。
  4. 可根据需要添加异常类型判断,比如只对网络类异常进行重试,其他异常直接进入断开状态。

内容的提问来源于stack exchange,提问作者osrl

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.14 10:04:56