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
相关产品推荐
相关产品推荐

