Swift Combine自定义Publisher如何在订阅后自动生成发送值
现有实现的核心问题是多订阅者场景下行为不符合预期:你使用了全局共享的PassthroughSubject,每接入一个新订阅者就调用一次startSending(),会导致已订阅的用户重复收到值,订阅序列混乱。同时直接在receive(subscriber:)中发送值也没有遵守Combine的背压处理规则,极端场景下可能出现丢值。
根据你实际的业务需求,有两种常见的实现方案:
场景1:每个新订阅者独立接收完整的1、2、3序列
这种场景不需要使用ConnectablePublisher,通过自定义Subscription实现即可,每个订阅者对应独立的发送逻辑,原生兼容背压规则:
class MyPublisher: Publisher { typealias Output = Int typealias Failure = Never func receive<S>(subscriber: S) where S : Subscriber, S.Failure == Failure, S.Input == Output { let subscription = MySubscription(subscriber: subscriber) subscriber.receive(subscription: subscription) } private class MySubscription<S: Subscriber>: Subscription where S.Input == Int, S.Failure == Never { private var subscriber: S? private let values = [1, 2, 3] private var sendIndex = 0 init(subscriber: S) { self.subscriber = subscriber } func request(_ demand: Subscribers.Demand) { var remainingDemand = demand while let subscriber = subscriber, sendIndex < values.count, remainingDemand > 0 { remainingDemand -= 1 let additionalDemand = subscriber.receive(values[sendIndex]) remainingDemand += additionalDemand sendIndex += 1 } if sendIndex >= values.count { subscriber?.receive(completion: .finished) } } func cancel() { subscriber = nil } } }
场景2:所有订阅者共享同一份输出序列,发送逻辑仅执行一次
如果你的发送逻辑成本较高(比如网络请求、批量计算),希望先挂载完所有订阅者后再统一触发一次发送,所有订阅者收到相同的序列,此时ConnectablePublisher正好匹配该需求。
你可以直接用Combine内置的makeConnectable()方法快速实现,不需要手动编写ConnectablePublisher的协议实现:
class MyPublisher: Publisher { typealias Output = Int typealias Failure = Never private let subject = PassthroughSubject<Int, Never>() func receive<S>(subscriber: S) where S : Subscriber, S.Failure == Failure, S.Input == Output { subject.receive(subscriber: subscriber) } // 仅在需要的时候主动调用触发发送 func startSending() { subject.send(1) subject.send(2) subject.send(3) subject.send(completion: .finished) } } // 使用示例 var cancellables = Set<AnyCancellable>() let publisher = MyPublisher() // 转换为ConnectablePublisher let connectablePublisher = publisher.makeConnectable() // 先挂载所有订阅者 connectablePublisher .sink { print("订阅者1收到:\($0)") } .store(in: &cancellables) connectablePublisher .sink { print("订阅者2收到:\($0)") } .store(in: &cancellables) // 所有订阅者挂载完成后,调用connect触发发送 connectablePublisher .connect() .store(in: &cancellables) // 也可以使用autoconnect()实现第一个订阅者挂载后自动触发发送 publisher .makeConnectable() .autoconnect() .sink { print("订阅者收到:\($0)") } .store(in: &cancellables)
如果需要第一个订阅者附着后自动调用startSending(),不需要手动调用connect(),直接在ConnectablePublisher后拼接.autoconnect()即可。
内容的提问来源于stack exchange,提问作者drekka
相关产品推荐
相关产品推荐

