如何手动完成Combine流?解决CurrentValueSubject导致流不完成问题
问题描述
我有一个整数流,需要通过使用CurrentValueSubject的结构体控制是否允许处理。由于CurrentValueSubject不会完成,若在flatMap()中执行处理,整体流(fooTheBars)仅当整数流为空时才会完成。
注释掉处理部分(.flatMap { bar in fooer.doFoo(with: bar) })后,因未涉及CurrentValueSubject,流可正常完成,但我需要无论项数多少,所有项处理完成后流都能正常完成。
当前输出
[0, 1, 2] doing foo with 0 doing foo with 1 doing foo with 2 received value received value received value
期望输出
[0, 1, 2] doing foo with 0 doing foo with 1 doing foo with 2 received value received value received value received completion // hooray, completion!
空流输出(符合预期)
[] received completion // hooray, completion!
注释flatMap()后的输出(正常完成)
[0, 1, 2] received value received value received value received completion // hooray, completion!
相关代码(注释上方不可修改)
struct Fooer { private let readyToFoo = CurrentValueSubject<Bool, Never>(true) func doFoo(with item: Int) -> AnyPublisher<Void, Never> { readyToFoo .filter { $0 } .map { _ in () } .handleEvents(receiveOutput: { print("doing foo with \(item)") // processing }) .delay(for: .seconds(1), scheduler: DispatchQueue.main) .eraseToAnyPublisher() } } func getBars() -> AnyPublisher<Int, Never> { let bars = (0..<(Bool.random() ? 0 : 3)) print(Array(bars)) return bars.publisher .eraseToAnyPublisher() } // can't change anything above this comment let fooer = Fooer() let fooTheBars = getBars() .flatMap { bar in fooer.doFoo(with: bar) } .sink(receiveCompletion: { completion in print("received completion") }, receiveValue: { _ in print("received value") })
解决方案
修改注释下方的代码
问题核心是doFoo返回的流依赖CurrentValueSubject(无限流),导致单个项的处理流永远不会完成,进而让flatMap无法触发整体流的完成。只需给每个doFoo的流添加.first()操作符,让它在发送一次输出后立即完成:
let fooer = Fooer() let fooTheBars = getBars() .flatMap { bar in fooer.doFoo(with: bar) .first() // 关键修改:取第一个输出后立即完成子流 } .sink(receiveCompletion: { completion in print("received completion // hooray, completion!") }, receiveValue: { _ in print("received value") })
这样每个子流都会在完成单次doFoo操作后发送完成信号,当getBars()的源流结束且所有子流都完成时,整个fooTheBars流就会触发完成回调。
关于CurrentValueSubject的设计模式评价
这种使用方式属于不良设计,原因如下:
doFoo的语义是"执行一次foo操作",但返回的却是无限流,违背了函数的语义预期——调用者理应得到代表单次操作完成的有限流,而非持续监听状态的流。- 无限流会破坏Combine的完成链语义,上游完成后下游无法正常结束,容易引发资源泄漏或逻辑错误。
更合理的设计方向:
- 如果用
readyToFoo控制操作权限,应让doFoo在状态允许时执行单次操作并完成,比如用.first()截断流,或者在状态满足时生成单次事件后结束。 - 若需要持续响应状态变化,应将状态监听逻辑与单次操作逻辑分离,避免语义混淆。
内容的提问来源于stack exchange,提问作者foxneSs
相关产品推荐
相关产品推荐

