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

如何为返回SharedFlow的挂起函数实现合规超时机制?

优化方案分析与实现

核心问题拆解

你当前实现的主要问题是:

  • 手动创建CoroutineScope(coroutineContext)会导致子协程生命周期与调用方协程脱节,调用方协程取消时超时逻辑仍会继续执行,存在资源泄漏风险
  • 未及时清理methodListeners中的监听器,长期运行会导致内存泄漏
  • 原MutableSharedFlow默认replay=0,调用方若延迟订阅会错过Loading状态

优化思路

  1. 用Flow操作符替代手动协程监听:通过merge合并业务流与超时流,更符合响应式编程范式,避免手动维护协程状态
  2. 绑定协程生命周期:使用当前协程上下文+SupervisorJob创建共享流的Scope,确保子协程随调用方协程取消而终止
  3. 自动清理资源:通过onCompletion操作符在流结束时自动移除监听器,避免内存泄漏
  4. 保证状态可达性:设置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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 13:52:53