如何使用Kotlin Flows编写请求响应代码并保证监听器已就绪
解决方案
1 架构层根治方案:统一缓存进站消息
从底层连接设计入手可以彻底解决时序问题,无需上层处理复杂的同步逻辑:
将连接收到的所有响应帧统一写入一个容量足够的Channel或带replay缓存的SharedFlow,所有上层接收操作都从这个缓存中读取。这样无论响应返回速度多快、上层接收逻辑启动多晚,消息都会先存在缓存中不会丢失。
示例实现:
// 全局进站消息缓存,replay设置为足够大的值避免消息溢出 val inboundBus = MutableSharedFlow<Any>(replay = 128) // 单独启动协程持续接收连接消息写入缓存 launch(Dispatchers.IO) { while (connection.isActive) { val msg = connection.receive<Any>() inboundBus.emit(msg) } } // 上层请求逻辑无需关心时序,直接过滤对应响应即可 suspend fun getServices(): GetServicesResponse { connection.send(GetServices(...)) return inboundBus.filterIsInstance<GetServicesResponse>().first() }
2 显式就绪同步方案(无需改造底层连接)
如果无法修改底层连接实现,可以用CompletableDeferred做显式的就绪通知,硬保证监听逻辑完全就绪后再发送请求。
2.1 单次请求响应场景
val receiveReady = CompletableDeferred<Unit>() val responseTask = async { // 第一行就发送就绪信号,保证后续receive执行前上层已经收到就绪通知 receiveReady.complete(Unit) connection.receive<GetServicesResponse>() } // 等待接收逻辑完全就绪后再发请求 receiveReady.await() connection.send(GetServices(...)) // 等待响应返回 val response = responseTask.await()
该方案不会依赖协程调度器的实现,只要receiveReady.await()返回就说明接收逻辑已经执行到即将调用receive的位置,不会出现响应丢失的问题。
2.2 Flow场景
同样用就绪信号的思路,将信号发送点放到Flow真正完成监听注册的位置即可:
// 扩展方法:给Flow添加就绪信号通知 fun <T> Flow<T>.withReadySignal(ready: CompletableDeferred<Unit>): Flow<T> = flow { ready.complete(Unit) emitAll(this@withReadySignal) } // 使用示例 val flowReady = CompletableDeferred<Unit>() val collectJob = launch { realFlow // 如果realFlow内部有嵌套协程,就把withReadySignal放到realFlow内部真正注册完监听的位置 .withReadySignal(flowReady) .collect { // 处理响应逻辑 } } // 等待Flow就绪后再发请求 flowReady.await() connection.send(GetServices(...))
如果realFlow内部存在嵌套协程,只需要将ready.complete(Unit)的调用位置移动到realFlow内部真正完成监听注册的代码之后即可,不受外层协程调度影响。
3 热流适配RxJava使用习惯
如果习惯RxJava订阅即就绪的逻辑,可以使用SharedFlow/StateFlow这类热流替代冷Flow:热流和RxJava的Observable逻辑一致,不存在冷流需要触发collect才启动的问题,订阅后会立即开始接收消息,不需要额外做同步处理。
内容的提问来源于stack exchange,提问作者David Stocking
相关产品推荐
相关产品推荐

