Swift Combine Publishers.CombineLatest是否存在背压处理Bug?
CombineLatest背压需求传播异常问题
问题描述
使用Publishers.CombineLatest过程中观测到需求管理相关异常:自定义实现.withLatestFrom运算符、通过限制上游需求做背压管理时,Publishers.CombineLatest仅向第二个上游发布者(latest2)传播下游需求,未向第一个上游发布者(latest1)传递对应需求,直接导致整条发布者链运行异常。
疑问点为该现象是否为CombineLatest的固有Bug,或是存在未被注意到的使用误区。
复现代码
func testCombineLatest() { let latest1 = PassthroughSubject<Int,Never>() let latest2 = PassthroughSubject<Int,Never>() var result:[[Int]] = [] var subscription:Subscription? let subscriber = AnySubscriber<(Int,Int),Never>( receiveSubscription: {sub in subscription = sub sub.request(.max(1)) }, receiveValue: { (v1,v2) in result.append([v1,v2]) return .max(1) }, receiveCompletion: {_ in} ) let publisher = Publishers.CombineLatest(latest1.print("Latest1"), latest2.print("Latest2")) .print("CombineLatest") publisher .subscribe(subscriber) latest1.send(1) latest2.send(1) latest1.send(2) //<- 该值发送后无任何效果 latest2.send(2) latest1.send(completion: .finished) latest2.send(completion: .finished) print("Result is:\(result)") XCTAssertEqual(result, [[1,1], [2,1], [2,2] ]) }
运行输出日志
Test Case '-[CAPTests.CAPTestPublisherExtensions testCombineLatest]' started. Latest1: receive subscription: (PassthroughSubject) Latest2: receive subscription: (PassthroughSubject) CombineLatest: receive subscription: (CombineLatest) CombineLatest: request max: (1) Latest1: request max: (1) Latest2: request max: (1) Latest1: receive value: (1) Latest2: receive value: (1) CombineLatest: receive value: ((1, 1)) CombineLatest: request max: (1) (synchronous) Latest2: request max: (1) (synchronous) Latest2: receive value: (2) CombineLatest: receive value: ((1, 2)) CombineLatest: request max: (1) (synchronous) Latest2: request max: (1) (synchronous) Latest1: receive finished Latest2: receive finished CombineLatest: receive finished Result is:[[1, 1], [1, 2]] error: -[CAPTests.CAPTestPublisherExtensions testCombineLatest] : XCTAssertEqual failed: ("[[1, 1], [1, 2]]") is not equal to ("[[1, 1], [2, 1], [2, 2]]")
解答
这是苹果Combine框架中Publishers.CombineLatest的固有实现Bug,不属于使用方式问题。
按照Combine的背压设计规范,当下游通过receiveValue返回新的.max(n)需求时,CombineLatest需要同时向两个上游发布者补充对应额度的请求,保证两个上游都有足够配额发送新值,才能正确完成“取两个上游最新值组合”的逻辑。
从运行日志可以直接定位问题:
- 初始订阅阶段,两个上游都正确收到了
.max(1)的初始需求 - 第一次组合值
(1,1)发送给下游、下游返回新的.max(1)需求后,只有Latest2收到了同步的.max(1)新请求,Latest1没有拿到任何新增的发送配额 - 此时调用
latest1.send(2),由于Latest1已经没有剩余的可发送额度,这个值会被PassthroughSubject直接丢弃,不会进入CombineLatest的内部值缓存,后续latest2.send(2)触发组合时,取到的还是Latest1缓存的旧值1,自然无法得到预期的[2,1]、[2,2]结果。
临时规避方案
- 对需要严格背压管控的链路,不要依赖系统
CombineLatest的默认背压传播逻辑,可以在两个上游发布者后拼接.buffer算子,配置足够的预取额度缓存上游值,避免值被无预期丢弃 - 自行实现符合背压规则的
CombineLatest算子,收到下游新需求时,同时向两个上游分配对应请求额度,不要仅向第二个上游传递需求。
内容的提问来源于stack exchange,提问作者Michael J
相关产品推荐
相关产品推荐

