源Publisher完成时订阅取消Subject的技术问题咨询
我创建了一个接收<Publisher>的类,内部用<PassthroughSubject>同时处理两个场景:
- 订阅传入的源 Publisher
- 手动调用
<send>方法发送值
示例代码(修正语法后)
import Combine // 假设你的 Scope 是封装了 AnyCancellable 数组的类型 typealias Scope = [AnyCancellable] class Adapter<T> { let innerSubject = PassthroughSubject<T, Never>() var scope = Scope() init(_ source: some Publisher<T, Never>) { source .subscribe(innerSubject) .store(in: &scope) innerSubject .sink(receiveValue: { debugPrint($0) }) .store(in: &scope) } func adapt(_ val: T) { innerSubject.send(val) } } func usage() { let adapter = Adapter<Int>(Empty()) // 改为 Empty(completeImmediately: false) 可临时解决 adapter.adapt(42) // 预期打印 42,但实际无输出 }
问题描述
当使用Empty()作为源 Publisher 时,内部的PassthroughSubject会因为Empty()默认立刻触发完成事件而被终止,导致后续手动调用adapt(_:)发送的值无法被接收。我预期的行为是仅取消源 Publisher 到 Subject 的订阅,而非让 Subject 本身终止(类似.NET中的行为),想请教:
- 是否遗漏了 Combine 中的某些机制?
- 是否需要用广播类的 Subject 来包装?
- 我原本以为 Publisher 支持多订阅者并天然实现多播,这个理解有问题吗?
问题根源
Combine 中的PassthroughSubject作为 Publisher,一旦接收到上游的finished事件,就会将自身标记为已完成状态:后续所有调用send(_:)发送的值都会被忽略,同时所有订阅者都会收到完成事件并取消订阅。这是 Subject 的核心行为,和.NET中对应的类型行为存在差异。
而你使用的Empty()默认参数是completeImmediately: true,会在订阅后立刻发送完成事件,直接终结了innerSubject。
解决方案
方案1:控制源 Publisher 的完成时机
正如你发现的,使用Empty(completeImmediately: false)可以避免源立刻发送完成事件,让innerSubject保持活跃。但这种方式依赖于源 Publisher 的具体实现,通用性不强。
方案2:拦截上游的完成事件
如果需要保留源 Publisher 的其他行为,但不想让它的完成事件传递给innerSubject,可以用handleEvents拦截并过滤完成事件:
init(_ source: some Publisher<T, Never>) { source .handleEvents(receiveCompletion: { completion in // 仅转发错误事件,忽略完成事件 if case .failure(let error) = completion { innerSubject.send(completion: .failure(error)) } }) .subscribe(innerSubject) .store(in: &scope) innerSubject .sink(receiveValue: { debugPrint($0) }) .store(in: &scope) }
这样上游的完成事件不会传递给innerSubject,但错误事件可以按需转发,保证innerSubject能继续接收手动发送的值。
方案3:用 Multicast 实现可靠多播
如果你需要让源 Publisher 的事件分发给多个订阅者,同时避免源的完成事件终结整个广播流,推荐使用Multicast操作符:
class Adapter<T> { let innerSubject = PassthroughSubject<T, Never>() var scope = Scope() init(_ source: some Publisher<T, Never>) { // 将源 Publisher 转为多播流,共享给所有订阅者 let multicasted = source.multicast(subject: innerSubject) // 订阅多播流的输出 multicasted .sink(receiveValue: { debugPrint($0) }) .store(in: &scope) // 连接源和多播 Subject,启动数据流 multicasted.connect() .store(in: &scope) } func adapt(_ val: T) { innerSubject.send(val) } }
Multicast的作用是让源 Publisher 仅被订阅一次,将事件分发给所有订阅者。关键在于,源的完成事件只会终结multicasted这个流,但innerSubject本身不会被终止(除非手动发送完成/失败事件),后续仍然可以通过adapt(_:)发送值。
关于多播的理解
Combine 中的普通 Publisher 默认是冷 Publisher,每次被订阅都会重新执行一次数据流逻辑。而多播操作(如multicast、share)是将冷 Publisher 转为热 Publisher,让多个订阅者共享同一个数据流。但即使是热 Publisher,一旦上游发送完成/失败事件,整个流还是会终止。如果需要一个永远不会自动终止的广播流,要么手动拦截上游的完成事件,要么使用PassthroughSubject这类不会自动终结的 Subject(只要不主动发送完成/失败事件,就会一直保持活跃)。
内容的提问来源于stack exchange,提问作者NSMustache

