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

Kotlin SharedFlow自定义onActive扩展时distinctUntilChanged()不生效如何解决

问题根本原因

  • Flow 默认是冷流类型,每新增一个订阅者都会独立触发一次上游流的采集逻辑。你当前代码中,isActiveFlow(包含map、distinctUntilChanged操作)没有做共享处理,所以每新增一个测试订阅者,都会单独订阅一次原始的subscriptionCount流,每个订阅者持有独立的distinctUntilChanged去重状态,自然就会出现多次重复触发isActive打印的问题。
  • 额外的叠加问题:你返回的流是通过flatMapLatest生成的,每次有新订阅者都会重新执行flatMapLatest的回调,当全局订阅数变化时,所有已存在的订阅者都会各自收到状态变更通知,进一步加剧了重复触发的问题。

修复方案

核心思路是把订阅计数的状态判断逻辑改成共享热流,让所有下游订阅者共用同一份上游订阅和去重状态,修改后的onActive扩展方法如下:

fun <T : Any> MutableSharedFlow<T>.onActive(
    block: suspend CoroutineScope.() -> Unit
): Flow<T> {
    val original = this
    // 将订阅计数判断逻辑转为共享热流,所有下游共用同一份状态
    val isActiveFlow: Flow<Boolean> = subscriptionCount
        .map {
            println("Class: Count is $it")
            it > 0
        }
        .distinctUntilChanged()
        // 关键:使用shareIn共享上游订阅,replay=1保证新订阅者直接拿到当前状态,跟随原始流的生命周期
        .shareIn(
            scope = CoroutineScope(Dispatchers.Unconfined),
            started = SharingStarted.WhileSubscribed(),
            replay = 1
        )

    return isActiveFlow.flatMapLatest { isActive ->
        println("Class: isActive is $isActive")
        if (isActive) {
            // 启动block协程,和当前流绑定,订阅数归0时flatMapLatest切换流会自动取消该协程
            coroutineScope {
                launch { block() }
                original
            }
        } else {
            original
        }
    }
}

修改后再运行测试,就能看到isActive只会在订阅数从0变正、以及从正变0的时候各打印一次,不会出现重复触发的问题。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.30 23:15:02