Swift Combine多线程问题:AnyPublisher与PassthroughSubject状态累积异常
多线程下Combine流合并与
scan状态累积异常问题 我在开发中通过合并多个Combine流维护状态,主流依赖scan(_:)操作符累积子流元素。其中一个子流由PassthroughSubject驱动,会被应用内多线程调用。
我清楚多线程向Subject发送事件时,Combine会自行调度线程,且scan(_:)本身并非线程安全,因此主流的状态累积会出现错误。原以为在流中使用receive(on:)能让状态累积正常工作,但仍观察到事件处理顺序混乱的情况。
已尝试以下方案均无效:
- 在监听Subject的子流中添加
receive(on:) - 为
receive(on:)操作符添加barrier属性 - 给Subject加锁(后续发现无意义,Subject自身已处理锁逻辑)
核心实现代码
struct MainStream { enum MainEvent { case start case setCount(Int) } struct MainState: Equatable { var count: Int } let queue: DispatchQueue let subInteractor: SubStream func build(_ upstream: AnyPublisher<MainEvent, Never>) -> AnyPublisher<MainState, Never> { let shared = upstream.share().eraseToAnyPublisher() let subStream = subInteractor .build( shared .map { _ in SubStream.SubEvent.observe } .eraseToAnyPublisher() ) .map { state in MainEvent.setCount(state.count) } return shared .merge(with: subStream) .scan(MainState(count: 0)) { accum, event in switch event { case .start: return accum case let .setCount(int): return MainState(count: int) } } .eraseToAnyPublisher() } } struct SubStream { struct SubState { let count: Int } enum SubEvent { case observe } let queue: DispatchQueue let subject: PassthroughSubject<Int, Never> func build(_ upstream: AnyPublisher<SubEvent, Never>) -> AnyPublisher<SubState, Never> { upstream .map { event in // 生产代码中会根据event处理,示例中无需关注 subject .receive(on: queue) .map { int in SubState(count: int) } } .switchToLatest() .receive(on: queue) .eraseToAnyPublisher() } }
多线程场景单元测试代码
var sut: MainStream! // ... 测试初始化代码省略 func test_multipleThreads() { let expectation0 = XCTestExpectation(description: #function) let expectation1 = XCTestExpectation(description: "expectation1") let expectation2 = XCTestExpectation(description: "expectation2") let expectation3 = XCTestExpectation(description: "expectation3") let queue1 = DispatchQueue(label: "queue1") let queue2 = DispatchQueue(label: "queue2") let queue3 = DispatchQueue(label: "queue3") let expectedResults = [0, 1, 2, 3] var receivedOutput: [Int] = [] sut.build(subject.eraseToAnyPublisher()) .collect(4) .sink { output in receivedOutput = output.map { $0.count } expectation0.fulfill() } .store(in: &cancellables) subject.send(.setCount(0)) queue1.async { [self] in interactorSubject.send(1) expectation1.fulfill() } queue2.async { [self] in interactorSubject.send(2) expectation2.fulfill() } queue3.async { [self] in interactorSubject.send(3) expectation3.fulfill() } wait(for: [expectation0, expectation1, expectation2, expectation3], timeout: 0.5) XCTAssertEqual(receivedOutput, expectedResults) }
该测试通过率仅约40%,希望理清线程逻辑的疏漏点并解决问题。
问题分析与解决方案
核心问题根源
PassthroughSubject的线程特性:PassthroughSubject内部通过锁保证send(_:)操作的原子性,但不保证下游接收事件的顺序与发送顺序完全一致。多线程发送事件时,Subject仅将事件加入队列,若下游使用并发队列调度,事件可能被乱序执行。scan的线程不安全:scan的累积闭包会在事件所在线程执行,若多个事件同时进入scan,会触发竞态条件,导致状态累积异常。receive(on:)使用时机错误:仅在子流中添加receive(on:)无法保证合并后的流在同一线程执行scan,主流事件仍可能在其他线程触发,破坏状态累积的原子性。
解决方案
确保所有进入scan的事件都在同一个串行队列上处理,保证状态累积的顺序和原子性,具体修改如下:
1. 在合并流后、scan前添加receive(on:)
修改MainStream的build方法,强制合并后的所有事件调度到指定串行队列:
func build(_ upstream: AnyPublisher<MainEvent, Never>) -> AnyPublisher<MainState, Never> { let shared = upstream.share().eraseToAnyPublisher() let subStream = subInteractor .build( shared .map { _ in SubStream.SubEvent.observe } .eraseToAnyPublisher() ) .map { state in MainEvent.setCount(state.count) } return shared .merge(with: subStream) // 关键:强制所有事件在同一个串行队列处理,保证scan的原子性与顺序 .receive(on: queue) .scan(MainState(count: 0)) { accum, event in switch event { case .start: return accum case let .setCount(int): return MainState(count: int) } } .eraseToAnyPublisher() }
注意:此处的queue必须是串行队列(默认DispatchQueue(label:)创建的即为串行队列),使用并发队列无法保证顺序。
2. 优化子流的receive(on:)使用
子流中switchToLatest()后的receive(on: queue)可移除,因为合并后的流已被调度到指定队列,避免重复调度:
func build(_ upstream: AnyPublisher<SubEvent, Never>) -> AnyPublisher<SubState, Never> { upstream .map { event in subject .receive(on: queue) .map { int in SubState(count: int) } } .switchToLatest() .eraseToAnyPublisher() }
3. 测试细节调整
多线程异步发送事件的执行顺序本身不确定,若需测试稳定通过,可调整测试逻辑:
- 若期望严格顺序,改用串行队列发送事件
- 若仅需验证最终状态包含所有预期值,修改断言逻辑为检查
receivedOutput是否包含所有expectedResults元素(不限制顺序)
为什么之前的方案无效?
- 子流添加
receive(on:)仅能保证子流事件在指定队列处理,但合并后的流仍可能在多线程执行scan,主流事件的调度线程不受控制。 barrier仅对并发队列有效,且要求所有操作都通过该队列提交才能保证顺序,与当前使用场景不匹配。- Subject自身的锁仅保证
send(_:)操作不并发执行,不保证下游处理的顺序。
内容的提问来源于stack exchange,提问作者d7774a59-f96c-44e8-8b2c-9b5a0e
相关产品推荐
相关产品推荐

