如何构建懒加载SharedFlow:延迟调用createFlow并镜像其返回流
要实现仅在第一个订阅者连接时才调用createFlow(),之后所有订阅者共享同一原流的需求,可以通过结合协程的延迟初始化特性与SharedFlow的转发能力来完成,以下是两种实用方案:
方案一:基于lazy+Deferred的代理SharedFlow
这种方式会创建一个代理SharedFlow,自动处理懒加载逻辑,且完全镜像原流的事件:
import kotlinx.coroutines.* import kotlinx.coroutines.flow.* // 替换为你的SharedFlow实际元素类型 typealias FlowElementType = String // 你的原挂起函数 suspend fun createFlow(): SharedFlow<FlowElementType> { // 示例hot流实现,实际替换为你的业务逻辑 return MutableSharedFlow<FlowElementType>().apply { println("createFlow 已执行(仅触发一次)") launch { repeat(5) { emit("原流事件 $it") delay(1000) } } }.asSharedFlow() } fun CoroutineScope.createLazyMirroredFlow(): SharedFlow<FlowElementType> { // 用lazy延迟初始化,确保createFlow仅被调用一次 val originalFlowDeferred by lazy { async { createFlow() } } // 创建转发原流事件的冷流 val mirrorFlow = flow { val originalFlow = originalFlowDeferred.await() emitAll(originalFlow) } // 转换为SharedFlow,配置订阅策略 return mirrorFlow.shareIn( scope = this, // 无订阅者时等待5秒再停止收集,避免频繁启停 started = SharingStarted.WhileSubscribed(5000), // 可根据原流的replay值调整,若原流replay固定可直接写死 replay = 0 ) } // 测试示例 fun main() = runBlocking { val scope = CoroutineScope(Dispatchers.Default) val lazyFlow = scope.createLazyMirroredFlow() println("第一个订阅者启动") val job1 = launch { lazyFlow.collect { println("订阅者1收到:$it") } } delay(2500) println("第二个订阅者启动") val job2 = launch { lazyFlow.collect { println("订阅者2收到:$it") } } delay(3000) job1.cancel() job2.cancel() scope.cancel() }
方案说明
- 懒加载触发:
originalFlowDeferred通过lazy初始化,只有当第一个订阅者触发流收集时,才会执行async { createFlow() },确保createFlow()仅调用一次。 - 事件转发:
mirrorFlow通过emitAll完全转发原流的所有事件,包括后续新增的事件。 - 订阅策略:
SharingStarted.WhileSubscribed(5000)保证有订阅者时持续收集原流事件,无订阅者时等待5秒再停止,避免不必要的资源消耗。
方案二:基于Mutex的直接缓存原流
如果不需要额外的代理流,可直接缓存原SharedFlow,让订阅者直接获取原流实例:
import kotlinx.coroutines.* import kotlinx.coroutines.flow.* import kotlinx.coroutines.sync.Mutex import kotlinx.coroutines.sync.withLock typealias FlowElementType = String suspend fun createFlow(): SharedFlow<FlowElementType> { // 原流实现 return MutableSharedFlow<FlowElementType>().apply { println("createFlow 已执行(仅触发一次)") launch { repeat(5) { emit("原流事件 $it") delay(1000) } } }.asSharedFlow() } class LazyFlowProvider(private val scope: CoroutineScope) { private val mutex = Mutex() private var cachedFlow: SharedFlow<FlowElementType>? = null suspend fun getLazyFlow(): SharedFlow<FlowElementType> { mutex.withLock { if (cachedFlow == null) { cachedFlow = createFlow() } return cachedFlow!! } } } // 测试示例 fun main() = runBlocking { val scope = CoroutineScope(Dispatchers.Default) val provider = LazyFlowProvider(scope) println("第一个订阅者启动") val job1 = launch { provider.getLazyFlow().collect { println("订阅者1收到:$it") } } delay(2500) println("第二个订阅者启动") val job2 = launch { provider.getLazyFlow().collect { println("订阅者2收到:$it") } } delay(3000) job1.cancel() job2.cancel() scope.cancel() }
方案说明
- 线程安全缓存:通过
Mutex保证多协程环境下,createFlow()仅被调用一次。 - 直接复用原流:订阅者获取的是原
SharedFlow实例,完全继承原流的所有特性(如replay缓存、溢出策略等)。
内容的提问来源于stack exchange,提问作者Alexey
相关产品推荐
相关产品推荐

