如何通过Combine的buffer操作符实现背压避免flatMap向上游请求无限需求
问题原因
- 最初写法的核心问题是操作符顺序错误:
map(some_preprocessing)放在buffer前面,Sequence类型的Range.Publisher默认会一次性推送全部100万个元素,导致所有预处理逻辑在buffer和flatMap的背压规则生效前就全部执行,直接占满内存。 - 你对
flatMap(maxPublishers: .max(32))的背压理解是正确的:它确实会在内部活跃发布者达到上限时,停止向上游请求新元素,但上游的无限制推送已经提前把所有元素处理完了,背压信号没有起到作用。 - 之前加buffer无效的原因同样是顺序问题:buffer在map下游,根本挡不住map提前处理所有元素。
正确实现方案
两种实现方式都可以达到「只有flatMap有空闲槽位时,才执行预处理、发起请求」的效果:
方案1:调整操作符顺序,将预处理放在buffer之后
直接把buffer放在序列发布者的紧邻下游,再执行预处理操作,保证背压信号能一路传到最上游的序列发布者,限制元素生成速度:
import Foundation import Combine let cancellable = (0..<1_000_000).publisher // 紧邻序列发布者加buffer,限制上游推送速度 .buffer(size: 32, prefetch: .byRequest, whenFull: .dropNewest) // 预处理放在buffer之后,只有下游请求时才会执行 .map(some_preprocessing) .flatMap(maxPublishers: .max(32)) { request in URLSession.dataTaskPublisher(for: request) .map(\.data) .catch { _ in Just(Data()) } } .sink { completion in print(completion) } receiveValue: { value in print(value) } // 替换固定时间的sleep,程序会一直运行直到所有请求处理完成 RunLoop.main.run()
方案2:预处理放到flatMap闭包内部(更稳妥)
直接把预处理逻辑移到flatMap的闭包里,只有当flatMap拿到空闲槽位、要处理当前元素时,才会执行预处理,完全避免多余的预处理调用:
import Foundation import Combine let cancellable = (0..<1_000_000).publisher .flatMap(maxPublishers: .max(32)) { index in // 预处理放到这里,只有要发起请求时才执行 let request = some_preprocessing(index) return URLSession.dataTaskPublisher(for: request) .map(\.data) .catch { _ in Just(Data()) } } .sink { completion in print(completion) } receiveValue: { value in print(value) } RunLoop.main.run()
参数说明
关于你之前困惑的buffer参数:
size:缓冲区最大存储的元素数量,这里设为和并发数一致的32即可prefetchStrategy: .byRequest:只有当下游主动请求新元素时,才会向上游要新的元素,不会提前预取多余元素whenFull:在prefetchStrategy: .byRequest的场景下,下游每次请求的元素数不会超过缓冲区大小,缓冲区永远不会满,所以这个参数的取值不会影响最终效果
内容的提问来源于stack exchange,提问作者Louis Lac
相关产品推荐
相关产品推荐

