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

RxSwift中如何让PublishSubject忽略错误并保持订阅不被销毁

RxSwift中如何让PublishSubject忽略错误保持订阅活跃

问题根因

首先明确Rx标准序列契约:任何Observable(包括PublishSubject)一旦触发onError或onCompleted事件,就会永久进入终止状态,后续所有发送给它的事件都会被直接丢弃,且所有已存在的订阅都会被自动销毁。

你现有代码不生效的原因:

  • catchError操作符的作用是捕获上游的错误事件,用你返回的新序列替代原上游继续向下游发送事件,但此时原pubSubj已经触发onError进入终止状态,你后续再给它发送next("3"),它本身已经不会处理了,自然下游收不到。
  • 你返回.never()只是让替代序列不会主动终止,但原上游已经不会产生任何新事件,所以没有后续输出。

解决方案

方案1:改造事件类型,从根源避免触发onError(最推荐)

不要直接调用Subject的onError方法,把错误也包装成普通next事件的一部分,这样Subject永远不会进入终止状态,订阅会一直存活。
示例代码:

// 元素类型改为Result,错误作为next事件的内容传输
let pubSubj = PublishSubject<Result<String, Error>>()
let obs = pubSubj
    .observeOn(ConcurrentDispatchQueueScheduler(qos: .background))

let dispose = obs.subscribe(onNext: { result in
    switch result {
    case .success(let str):
        print("***** received onNext: \(str)\n")
    case .failure(let err):
        print("***** catch error \(err)\n")
        // 错误在这里单独处理,不会影响整个流的存活
    }
}, onCompleted: {
    print("***** completed\n")
}, onDisposed: {
    print("***** onDisposed\n")
})
                    
pubSubj.on(.next(.success("2")))
// 错误作为next事件发送,不会终止Subject
pubSubj.on(.next(.failure(RxError.overflow)))
            
pubSubj.on(.next(.success("3"))) // 可正常接收

该方案符合Rx规范,无额外订阅开销,是生产环境首选。

方案2:兼容onError调用,自动重订阅

如果因为外部限制无法修改Subject的事件类型,必须兼容onError调用,可以通过retry操作符实现自动重订阅,注意每次重订阅需要生成新的Subject实例(旧实例已经终止无法复用)。
示例代码:

// 维护当前可用的Subject引用
var currentPubSubj: PublishSubject<String>?
let obs = Observable<String>.deferred {
    let newSubj = PublishSubject<String>()
    currentPubSubj = newSubj
    return newSubj
}
.observeOn(ConcurrentDispatchQueueScheduler(qos: .background))
.catch { err in
    print("***** catchError \(err)\n")
    return .empty() // 结束本次序列触发重试
}
.retry() // 错误发生后自动重新订阅,生成新的Subject实例

let dispose = obs.subscribe(onNext: { str in
    print("***** received onNext: \(str)\n")
}, onError: { err in
    print("***** received onErr: \(err)\n")
}, onCompleted: {
    print("***** completed\n")
}, onDisposed: {
    print("***** onDisposed\n")
})

currentPubSubj?.on(.next("2"))
currentPubSubj?.onError(RxError.overflow) // 触发错误自动重订阅

// 注意:需要获取新生成的currentPubSubj实例发送后续事件
currentPubSubj?.on(.next("3"))

该方案缺点是错误发生前未处理的事件会丢失,且需要维护当前可用的Subject实例,仅作为备选。

补充说明

RxSwift中catchError的功能已经对齐其他Rx实现的onErrorResumeNext,你之前的用法未生效不是操作符的问题,是对Subject的终止特性认知不足导致的。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 03:24:04