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

RxSwift事件重复触发两次问题排查与解决求助

RxSwift中事件重复触发的问题排查与解决

问题描述

我搭建了一个结合内部与外部Action的沙盒代码,尽可能简化以复现事件重复触发的问题。

初始示例代码

import PlaygroundSupport
import RxSwift

class Sandbox {
    let publisher: PublishSubject<Int>

    private let disposeBag = DisposeBag()

    init() {
        self.publisher = PublishSubject()

        let publisher2 = publisher.debug()
        let publisher3 = PublishSubject<Int>()
        let publisherMerge = PublishSubject.merge([publisher2, publisher3])

        let operation = publisherMerge
            .map { $0 + 2 }

        let operation2 = operation
            .map { $0 + 3 }

        // 沙盒中永远不会触发
        operation
            .filter { $0 < 0 }
            .flatMapLatest(Sandbox.doSomething)
            .subscribe { print("Operation ", $0) }
            .disposed(by: disposeBag)

        operation2
            .subscribe { print("Operation2 ", $0) }
            .disposed(by: disposeBag)
    }

    static func doSomething(value: Int) -> Observable<Int> {
        return .just(value)
    }
}

let disposeBag = DisposeBag()
let sandbox = Sandbox()

Observable<Int>
    .interval(.seconds(3), scheduler: MainScheduler.instance)
    .subscribe(onNext: { _ in sandbox.publisher.onNext(0) })
    .disposed(by: disposeBag)

PlaygroundPage.current.needsIndefiniteExecution = true

初始运行日志

2022-11-07 21:16:15.292: Sandbox.playground:12 (init()) -> subscribed
2022-11-07 21:16:15.293: Sandbox.playground:12 (init()) -> subscribed
2022-11-07 21:16:18.301: Sandbox.playground:12 (init()) -> Event next(0)
2022-11-07 21:16:18.302: Sandbox.playground:12 (init()) -> Event next(0)
Operation2  next(5)

可以看到事件被触发了两次,现需排查问题原因并找到确保事件仅触发一次的方法。


补充示例代码(贴近实际场景)

为更清晰展示需求,补充如下示例代码:

import Foundation
import PlaygroundSupport
import RxSwift

struct ViewModel: Equatable {
    enum State: Equatable {
        case initialized
    }
    
    let state: State
}

enum Reducer {
    enum Action {
        case dummyReducerAction(String)
    }

    enum Effect {
        case dummyEffect
    }

    struct State: Equatable {
        let viewModel: ViewModel
        let effect: Effect?
    }

    static func reduce(state: State, action: Action) -> State {
        let viewModel: ViewModel = state.viewModel
        let _: Effect? = state.effect
        let viewModelState: ViewModel.State = viewModel.state
        let noChange = State(viewModel: viewModel, effect: nil)

        switch (action, viewModelState) {
        case (.dummyReducerAction(let s), _):
            print(s, Date())
            return noChange
        }
    }

    static func fromAction(_ action: Interactor.Action) -> Action {
        switch action {
        case .dummyAction:
            return .dummyReducerAction("fromAction")
        }
    }

    static func performAsyncSideEffect(effect: Effect) -> Observable<Action> {
        switch effect {
        case .dummyEffect:
            return .just(.dummyReducerAction("performSideEffect"))
        }
    }
}

class Interactor {
    enum Action {
        case dummyAction
    }

    let action: PublishSubject<Action>
    let viewModel: BehaviorSubject<ViewModel>

    private let disposeBag = DisposeBag()

    init() {
        self.action = PublishSubject()
        
        let initialViewModel = ViewModel(state: .initialized)
        let initialReducerState = Reducer.State(viewModel: initialViewModel, effect: nil)
        
        let externalAction = action.map(Reducer.fromAction).debug()
        let internalAction = PublishSubject<Reducer.Action>()
        let allActions = PublishSubject.merge([externalAction, internalAction])
        
        self.viewModel = BehaviorSubject(value: initialViewModel)
        
        let reducerState = allActions
            .scan(initialReducerState, accumulator: Reducer.reduce)
        
        let viewModelObservable = reducerState
            .map { $0.viewModel }
            .distinctUntilChanged()
        
        reducerState
            .compactMap { $0.effect }
            .flatMapLatest(Reducer.performAsyncSideEffect)
            .subscribe(internalAction)
            .disposed(by: disposeBag)
        
        viewModelObservable
            .subscribe(onNext: viewModel.onNext)
            .disposed(by: disposeBag)
    }
}

let disposeBag = DisposeBag()
let interactor = Interactor()

Observable<Int>
    .interval(.seconds(5), scheduler: MainScheduler.instance)
    .subscribe(onNext: { _ in interactor.action.onNext(.dummyAction) })
    .disposed(by: disposeBag)

PlaygroundPage.current.needsIndefiniteExecution = true

添加share()后的运行日志

2022-11-08 11:46:05.298: Sandbox.playground:71 (init()) -> subscribed
2022-11-08 11:46:10.310: Sandbox.playground:71 (init()) -> Event next(dummyReducerAction("fromAction"))
fromAction 2022-11-08 02:46:10 +0000
fromAction 2022-11-08 02:46:10 +0000

未添加share()的运行日志

2022-11-08 11:56:33.369: Sandbox.playground:71 (init()) -> subscribed
2022-11-08 11:56:33.370: Sandbox.playground:71 (init()) -> subscribed
2022-11-08 11:56:38.383: Sandbox.playground:71 (init()) -> Event next(dummyReducerAction("fromAction"))
fromAction 2022-11-08 02:56:38 +0000
2022-11-08 11:56:38.385: Sandbox.playground:71 (init()) -> Event next(dummyReducerAction("fromAction"))
fromAction 2022-11-08 02:56:38 +0000

即使添加了share()操作符,Reducer仍会触发两次并打印重复内容。


问题原因

核心原因在于RxSwift中普通Observable是冷序列,每次订阅都会重新执行序列中的操作:

  1. 初始示例中:publisherMerge被operation和operation2两次订阅,导致上游的publisher2(即publisher.debug())被订阅两次,因此每次事件都会触发两次debug输出。
  2. 补充示例中:allActions被reducerState订阅,而reducerState又被两个下游(处理effect的订阅、处理viewModel的订阅)订阅,这导致allActions对应的序列被执行两次,进而让Reducer.reduce被调用两次。

直接添加默认share()无效的原因是:默认share(replay: 0, scope: .whileConnected)仅在有活跃订阅时共享序列,若两个订阅几乎同时发起(比如在init中连续订阅),仍可能导致序列被重新订阅一次。


解决方案

方案1:使用share(replay: 0, scope: .forever)

将reducerState(或上游的allActions)转换为永久共享序列,确保所有订阅共享同一个序列执行:

let reducerState = allActions
    .scan(initialReducerState, accumulator: Reducer.reduce)
    .share(replay: 0, scope: .forever)

方案2:使用PublishSubject作为中间转发

将reducerState的结果转发到一个Subject中,所有下游订阅这个Subject:

let reducerStateSubject = PublishSubject<Reducer.State>()
allActions
    .scan(initialReducerState, accumulator: Reducer.reduce)
    .subscribe(reducerStateSubject)
    .disposed(by: disposeBag)

let viewModelObservable = reducerStateSubject
    .map { $0.viewModel }
    .distinctUntilChanged()

reducerStateSubject
    .compactMap { $0.effect }
    .flatMapLatest(Reducer.performAsyncSideEffect)
    .subscribe(internalAction)
    .disposed(by: disposeBag)

方案3:合并下游订阅

如果两个下游订阅逻辑可以合并,尽量减少订阅次数。比如在初始示例中,让operation2基于operation的共享序列,或者合并两个订阅的处理逻辑。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 18:55:13