多线程并发调用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
相关产品推荐
相关产品推荐

