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

Swift Combine异步环境下发布事件丢失问题排查

代码问题分析:订阅消息数量不稳定的原因

问题代码

func synchronized<T>(_ lock: NSLock, closure:() throws -> T) rethrows -> T {
    lock.lock()
    let r = try closure()
    lock.unlock()
    return r
}

@objcMembers class Main : NSObject {
    static var lock = NSLock()
    static var subscribers: [Int: AnyCancellable] = [:]
    static var signal = DispatchSemaphore(value: 0)
    static func perform() {
        for i in 0 ..< 5 {
            Task {
                await subscribe(i)
            }
        }
        Task {
            try! await Task.sleep(for:.seconds(1))
            BackgroundNetwork.shared.results.append(1)
        }
        signal.wait()
    }
}

// 模拟网络请求:发起请求并等待响应
func subscribe(_ i: Int) async {
    await withCheckedContinuation { ctx in
        let subscriber = BackgroundNetwork.shared.dataSource().first().sink { newValue in
            print("Pos\(i). \(newValue)")
            ctx.resume()
        }
        synchronized(Main.lock) {
            Main.subscribers[i] = subscriber
        }
    }
}

class BackgroundNetwork : ObservableObject {
    static var shared = BackgroundNetwork()
    @Published var results: [Int] = []
    
    func dataSource() -> AnyPublisher<[Int], Never> {
        return $results
            .dropFirst() // 避免发布当前已有值
            .eraseToAnyPublisher()
    }
}

核心问题

消息数量波动的根源不是锁的使用问题,而是订阅者生命周期管理和任务调度的竞态条件:

  1. 订阅者持有时机滞后
    subscribe函数中,创建subscriber后先绑定了订阅流,再通过锁存入Main.subscribers。而AnyCancellable的特性是:一旦实例没有被任何对象持有,就会被销毁,订阅也会自动取消。如果在subscriber还没存入字典前,results.append(1)就触发了发布,这个未被持有的subscriber会立即销毁,对应的订阅失效,自然收不到消息。

  2. 固定等待不可靠
    用Task.sleep(1秒)等待订阅创建的逻辑不严谨——系统对Task的调度是不确定的,1秒内可能有部分订阅任务还没执行到存储subscriber的步骤,此时触发发布,这部分未完成的订阅就会丢失。

修复方案

1. 确保订阅者被立即持有

调整subscribe逻辑,创建subscriber后立刻加锁存入字典,避免因未持有导致的提前取消:

func subscribe(_ i: Int) async {
    await withCheckedContinuation { ctx in
        let subscriber = BackgroundNetwork.shared.dataSource().first().sink { newValue in
            print("Pos\(i). \(newValue)")
            ctx.resume()
            // 订阅完成后移除,避免内存泄漏
            synchronized(Main.lock) {
                Main.subscribers.removeValue(forKey: i)
            }
        }
        // 立即存储订阅者,确保持有关系生效
        synchronized(Main.lock) {
            Main.subscribers[i] = subscriber
        }
    }
}

2. 等待所有订阅完成后再触发发布

替换固定时间等待,改为等待所有订阅任务执行完成,确保所有订阅都已生效:

static func perform() {
    // 启动所有订阅任务并等待全部完成
    async let task0 = subscribe(0)
    async let task1 = subscribe(1)
    async let task2 = subscribe(2)
    async let task3 = subscribe(3)
    async let task4 = subscribe(4)
    
    // 等待所有订阅初始化完成
    _ = await [task0, task1, task2, task3, task4]
    
    // 此时所有订阅都已生效,触发发布
    BackgroundNetwork.shared.results.append(1)
    
    // 给回调留处理时间,再结束信号
    Task {
        try! await Task.sleep(for: .seconds(0.1))
        signal.signal()
    }
}

补充说明

NSLock的使用是正确的,它确实保护了subscribers字典的线程安全。问题的本质是没有保证订阅生效时机早于发布触发时机,通过上述调整,就能确保5个订阅都能稳定收到发布消息。

内容的提问来源于stack exchange,提问作者NoobOfCpp

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 14:47:02