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

Kotlin RX的zipWith函数体是执行一次还是每个观察者各执行一次?

问题原因

这是RxJava冷流的默认特性导致的:setupReadListener()返回的Flowable默认属于冷流,每新增一个订阅者都会重新触发整条上游链路的执行:

  • 每个订阅者会独立订阅上游的readStartPublishProcessor和dataFlowable
  • 每个订阅者都会单独执行zipWith的合并函数,所以执行次数和订阅者数量完全一致

Rx中所有转换操作符默认返回的都是冷流,只有手动配置共享策略,才会实现上游执行结果多订阅者复用的效果。

解决方案

在最终返回的Flowable末尾添加共享操作符即可,根据业务场景可以选择以下两种方案:

方案1:使用.share()(适用绝大多数场景)

share操作符会让多个订阅者共享同一条上游链路,上游所有操作(包括zipWith的合并函数)只会执行一次,结果会下发给所有当前已订阅的观察者,修改后的代码如下:

private fun setupReadListener(): Flowable<ReadEvent> {
    val dataFlowable = dataPacketPublishProcessor.buffer(readDonePublishProcessor)
    
    return readStartPublishProcessor
        .zipWith(other = dataFlowable) { readStart, dataPackets ->
            Log.d(tag, "ReadEvent with ${dataPackets.size} data packets")
            ReadEvent(event = readStart, payload = combinePackets(dataPackets))
        }
        .share() // 新增此行实现多订阅者共享上游结果
}

share采用引用计数机制,当所有订阅者都取消订阅后,上游会被自动销毁,新的订阅者接入时会重新触发上游执行。

方案2:使用.replay(1).refCount()(需要缓存最近一次结果的场景)

如果你希望新订阅者接入时可以收到最近已经发射过的一次ReadEvent,可以用这个组合,它会在共享上游的同时缓存最近1次的发射结果下发给新订阅者。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 04:06:03