Combine使用share()后多filter流合并缺失输出的原因及处理方案
Combine 中的 share() 运算符会将上游冷发布者转换为热多播发布者,具有两个核心特性:
- 仅在第一次收到订阅时触发一次上游订阅,后续所有新订阅者共享同一条上游事件流
- 无历史事件缓存能力,上游一旦发送完成/错误事件,后续新的订阅者会直接收到结束事件,不会拿到任何历史值
你给出的示例中,Merge 操作符会按顺序订阅上游的两个过滤发布者:先订阅 filter1,此时触发 share() 上游的数组发布者立刻发送全部1、2、3、4、5值并发送完成事件;等 Merge 再订阅 filter2 时,上游流已经结束,自然不会收到2这个值。
你有 RxSwift 使用经验的话可以对应参考:Combine 的 share() 等价于 RxSwift 中 share(replay: 0, scope: .forever) 的行为。
至于移除 share() 后运行正常,是因为此时每个过滤发布者都会单独订阅上游的数组发布者,相当于两次独立的序列遍历,所以都能拿到完整的输入值,但这种方案放到API请求场景下就会触发两次重复网络调用,不符合你的需求。
要实现「仅发起一次上游请求,所有下游过滤逻辑都能拿到完整事件流」的需求,可选择以下方案:
方案1:手动控制多播连接时机
使用 makeConnectable() 将多播发布者的启动时机交由你手动控制,等所有下游订阅都完成后再触发上游事件发送,是最通用的解决方案:
import Combine let publisher = [1, 2, 3, 4, 5] .publisher .share() .makeConnectable() // 禁止自动触发上游订阅 let filter1 = publisher .filter { $0 == 1 } .print("filter1") let filter2 = publisher .filter { $0 == 2 } .print("filter2") let cancellable = Publishers .Merge(filter1, filter2) .sink { print("Result is: \($0)") } // 所有订阅准备完毕后,手动触发上游事件发送 publisher.connect()
运行后即可正常得到 1 和 2 两个输出结果,且上游仅执行一次。
方案2:适配API请求场景的优化方案
如果你的上游是仅返回单次结果的API请求,可以直接用 Future 封装请求逻辑,Future 本身自带共享执行结果的能力,无需额外添加 share(),所有订阅者都会拿到同一份请求结果,不会发起重复请求。
方案3:需要支持晚到订阅者的场景
如果业务存在「上游已经发送完结果,后续新的订阅者也需要拿到历史请求结果」的需求,可以自己实现带重播能力的多播发布者,或使用 CurrentValueSubject 作为多播的中间层,缓存最近的1次请求结果供后续订阅者使用。
内容的提问来源于stack exchange,提问作者skornos

