如何为返回SharedFlow的挂起函数实现合规超时机制?
优化方案分析与实现
核心问题拆解
你当前实现的主要问题是:
- 手动创建
CoroutineScope(coroutineContext)会导致子协程生命周期与调用方协程脱节,调用方协程取消时超时逻辑仍会继续执行,存在资源泄漏风险 - 未及时清理
methodListeners中的监听器,长期运行会导致内存泄漏 - 原
MutableSharedFlow默认replay=0,调用方若延迟订阅会错过Loading状态
优化思路
- 用Flow操作符替代手动协程监听:通过
merge合并业务流与超时流,更符合响应式编程范式,避免手动维护协程状态 - 绑定协程生命周期:使用当前协程上下文+
SupervisorJob创建共享流的Scope,确保子协程随调用方协程取消而终止 - 自动清理资源:通过
onCompletion操作符在流结束时自动移除监听器,避免内存泄漏 - 保证状态可达性:设置
replay=1确保订阅者能收到最新状态(包括Loading)
完整优化代码
import kotlinx.coroutines.* import kotlinx.coroutines.flow.* import java.util.* import java.util.concurrent.ConcurrentHashMap class Client { // 客户端独立协程Scope,处理消息接收后的挂起调用 private val clientScope = CoroutineScope(Dispatchers.IO + SupervisorJob()) // 用ConcurrentHashMap保证多线程环境下的安全访问 private val methodListeners = ConcurrentHashMap<String, MyCallListener>() sealed interface State { object Loading : State data class Success(val message: String) : State object Error : State } interface MyCallListener { suspend fun onResult(result: State) } suspend fun call(): SharedFlow<State> { val id = UUID.randomUUID().toString() // 初始化MutableSharedFlow,设置replay=1确保订阅者能收到Loading状态 val mutableFlow = MutableSharedFlow<State>(replay = 1) // 注册监听器 methodListeners[id] = object : MyCallListener { override suspend fun onResult(result: State) { mutableFlow.emit(result) } } // 发送Loading状态(replay=1保证emit不会失败) mutableFlow.emit(State.Loading) // 发送WebSocket消息(挂起函数,确保消息发送完成后再继续) send(methodMessage) // 构建超时流:5秒后发送Error状态 val timeoutFlow = flow { delay(5000) emit(State.Error) } // 合并业务流与超时流,处理状态去重和资源清理 return merge(mutableFlow, timeoutFlow) .distinctUntilChanged() // 避免重复发送相同状态 .onCompletion { // 无论流正常完成还是被取消,都清理监听器 methodListeners.remove(id) } .shareIn( // 创建绑定当前协程上下文的Scope,子协程随调用方协程终止 scope = CoroutineScope(coroutineContext + SupervisorJob(coroutineContext[Job])), // 订阅者取消后延迟5秒停止共享,避免频繁重建流 started = SharingStarted.WhileSubscribed(5000), replay = 1 ) } // 处理WebSocket消息接收,用客户端Scope启动协程调用挂起函数 fun onReceiveMessage(id: String, message: String) { val listener = methodListeners[id] ?: return clientScope.launch { listener.onResult(State.Success(message)) } } // 客户端销毁时必须调用,取消所有协程避免泄漏 fun destroy() { clientScope.cancel() methodListeners.clear() } // 模拟WebSocket发送挂起函数 private suspend fun send(message: String) { // 实际业务逻辑:发送WebSocket消息 } // 模拟methodMessage private val methodMessage = "test-method" } class MyUseCase( val client: Client ) { suspend fun foo() = client.call() }
关键优化点说明
- 协程生命周期绑定:
CoroutineScope(coroutineContext + SupervisorJob(coroutineContext[Job]))确保超时逻辑的协程会随调用方协程的取消而终止,避免无用的资源占用 - 资源自动清理:
onCompletion操作符会在流结束(包括正常完成、超时、取消)时触发,自动移除监听器,解决内存泄漏问题 - 状态可靠性:
replay=1和SharingStarted.WhileSubscribed()保证订阅者不会错过初始的Loading状态,同时在无订阅者时自动释放资源 - 线程安全:
ConcurrentHashMap替代普通MutableMap,避免多线程环境下的并发修改异常 - 挂起函数调用修正:客户端内部维护独立的
clientScope,解决在非挂起函数onReceiveMessage中调用挂起函数onResult的编译问题
内容的提问来源于stack exchange,提问作者osrl
相关产品推荐
相关产品推荐

