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

多线程并发调用getStream与update时,Combine转AsyncStream方案是否线程安全?

线程安全性分析与修复方案

你的这个实现不具备完整的线程安全性,核心问题出在cancellables集合的并发访问上,具体分析和修复方案如下:

线程安全风险点分析

  • PassthroughSubject的send(_:)方法本身是线程安全的,Combine框架原生支持在任意线程调用该方法发送事件,所以update(value:)方法的多线程调用是安全的。
  • 问题出在getStream()方法中对cancellables集合的操作:Set<AnyCancellable>不是线程安全的数据结构,当多个线程同时调用getStream()时,会并发执行.store(in: &cancellables)操作——这本质是对Set的插入操作,并发修改非线程安全的集合会触发数据竞争,甚至导致程序崩溃。
  • 补充:AsyncStream.Continuation的yield(_:)和finish()方法是线程安全的,所以从Subject接收值后调用continuation的部分不存在线程安全问题。

修复方案

要解决这个问题,我们需要给cancellables的访问加上线程同步保护,最常用的方式是用串行DispatchQueue来包裹所有对cancellables的修改操作:

final class SomeNotifier {
    static let shared = SomeNotifier()
    
    private let value = PassthroughSubject<String, Never>()
    private var cancellables: Set<AnyCancellable> = []
    // 创建串行队列用于保护cancellables的并发访问
    private let cancellableQueue = DispatchQueue(label: "com.yourapp.SomeNotifier.cancellableQueue")
    
    private init() {}
    
    func getStream() -> AsyncStream<String> {
        return AsyncStream { [weak self] continuation in
            guard let self = self else { return }
            
            let cancellable = self.value.sink { completion in
                switch completion {
                case .finished, .failure:
                    continuation.finish()
                }
            } receiveValue: { value in
                continuation.yield(value)
            }
            
            // 在串行队列中执行store操作,避免并发修改
            self.cancellableQueue.sync {
                self.cancellables.insert(cancellable)
            }
            
            // 当continuation被取消时,自动移除对应的cancellable,避免内存泄漏
            continuation.onTermination = { [weak self] _ in
                self?.cancellableQueue.sync {
                    self?.cancellables.remove(cancellable)
                }
            }
        }
    }
    
    func update(value: String) {
        self.value.send(value)
    }
}

额外优化说明

修复后的代码添加了continuation.onTermination的处理:当AsyncStream的迭代器被取消(比如for-await-in循环退出)时,自动从cancellables中移除对应的订阅,避免内存泄漏——这是原实现中缺失的重要逻辑,原实现会导致订阅者即使不再监听,cancellable也一直留在集合中,造成不必要的内存占用。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 21:12:16