You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.04 20:18:03