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

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%,希望理清线程逻辑的疏漏点并解决问题。


问题分析与解决方案

核心问题根源

  1. PassthroughSubject的线程特性:PassthroughSubject内部通过锁保证send(_:)操作的原子性,但不保证下游接收事件的顺序与发送顺序完全一致。多线程发送事件时,Subject仅将事件加入队列,若下游使用并发队列调度,事件可能被乱序执行。
  2. scan的线程不安全:scan的累积闭包会在事件所在线程执行,若多个事件同时进入scan,会触发竞态条件,导致状态累积异常。
  3. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 00:28:10