RxJava结合observeOn时startWith操作符失效问题咨询
问题根源:
observeOn的错误处理逻辑导致事件丢弃 这个问题的核心在于observeOn操作符默认的错误处理行为——当它收到上游的onError事件时,会直接清空待处理的事件队列,优先传递错误,导致startWith发射的1被丢弃了。
先解释两种场景的差异:
1. 无observeOn的同步流
当你没有添加observeOn时,整个流是同步执行的:
- 订阅触发后,
startWith立即发射1,下游订阅者瞬间收到onNext并输出内容; - 随后
startWith订阅源Observable.error,该Observable立即发射onError,下游捕获异常。
这个流程是线性同步的,所以符合你的预期。
2. 添加observeOn后的异步调度流
当加入observeOn(AndroidSchedulers.mainThread())后,事件的传递逻辑变了:
- 订阅触发后,
startWith发射1,observeOn会把这个onNext事件放到主线程的消息队列中,等待主线程的Looper循环处理; - 紧接着,
startWith订阅源Observable.error,这个Observable会在当前订阅线程(也就是调用subscribe的线程)立即发射onError事件; observeOn收到onError后,会严格遵循RxJava的流终止规则:立即清空队列里还没处理的onNext事件(就是那个1),直接把onError传递给下游订阅者。
这就导致你看不到onNext的输出,只收到了异常。
解决方案:开启delayError参数
observeOn有一个重载方法,可以传入delayError参数,让错误事件延迟到所有正常事件处理完成后再传递:
Observable.error<Int>(Exception()) .startWith(1) .observeOn(AndroidSchedulers.mainThread(), true) // 第二个参数设为true,开启延迟错误 .subscribe( { println("Item is $it") }, { it.printStackTrace() } )
当delayError = true时,observeOn会把收到的错误事件缓存起来,等队列里所有的onNext事件都处理完毕后,再把错误传递给下游。这样就能先输出Item is 1,随后再抛出异常,和你最初的预期一致。
如果你的场景需要更复杂的错误处理(比如捕获错误后继续流的执行),还可以考虑onErrorResumeNext、onErrorReturn等操作符,但针对当前问题,delayError是最直接的解决方式。
内容的提问来源于stack exchange,提问作者O10
相关产品推荐
相关产品推荐

