如何在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) }
逻辑说明
- 外层
while(true)让流保持活跃,除非WebSocket正常关闭或遇到无法处理的异常。 - 自动重试阶段:连接失败后,最多尝试3次自动重连,每次间隔500ms。
- 手动重试阶段:3次自动重试失败后,流进入挂起状态,等待
retryTrigger的信号;调用retryWith时发送信号,流会立即重新尝试建立连接。 - 可根据需要添加异常类型判断,比如只对网络类异常进行重试,其他异常直接进入断开状态。
内容的提问来源于stack exchange,提问作者osrl
相关产品推荐
相关产品推荐

