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

如何构建懒加载SharedFlow:延迟调用createFlow并镜像其返回流

实现懒加载并镜像SharedFlow的方案

要实现仅在第一个订阅者连接时才调用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()
}

方案说明

  1. 懒加载触发:originalFlowDeferred通过lazy初始化,只有当第一个订阅者触发流收集时,才会执行async { createFlow() },确保createFlow()仅调用一次。
  2. 事件转发:mirrorFlow通过emitAll完全转发原流的所有事件,包括后续新增的事件。
  3. 订阅策略: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()
}

方案说明

  1. 线程安全缓存:通过Mutex保证多协程环境下,createFlow()仅被调用一次。
  2. 直接复用原流:订阅者获取的是原SharedFlow实例,完全继承原流的所有特性(如replay缓存、溢出策略等)。

内容的提问来源于stack exchange,提问作者Alexey

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 05:25:56