RxJava中Flowable cache转Single出现死锁问题求解
解析RxJava缓存流阻塞问题:singleOrError/singleElement vs firstElement
这是个很有意思的RxJava线程阻塞问题,咱们一步步拆解背后的原因:
先理清代码的执行流程
你的代码里,cachedFlowable是Flowable.just(1).cache(),当调用blockingSubscribe()时,main线程会触发整个流的执行:
- 源
Flowable.just(1)发射元素1,进入doOnNext回调,打印doOnNext 1。 - 在
doOnNext里,你又订阅了cachedFlowable并调用blockingGet(),这里的关键差异来自操作符的行为和cache()的内部状态。
先搞懂cache()的核心状态
cache()的逻辑是:当第一个订阅者出现时,它会订阅源Flowable,缓存所有发射的元素,同时转发给订阅者。而Flowable.just(1)是同步发射的,它的执行顺序是:
- 发射
1→ 执行所有订阅者的onNext(也就是你的doOnNext) → 等onNext完全执行完毕 → 发送onComplete信号。
逐个分析操作符的行为
1. singleOrError()和singleElement()为什么阻塞?
这两个操作符的核心约束是:流必须恰好发射一个元素,并且必须收到onComplete信号来确认没有更多元素。
当你在doOnNext里调用cachedFlowable.singleOrError().blockingGet()时:
cachedFlowable已经缓存了1,所以新订阅会立即收到这个元素,但singleOrError()还在等待onComplete信号——它需要确认不会有第二个元素发射。- 但源的
onComplete什么时候发送?要等第一个订阅的doOnNext执行完才行!而main线程现在正卡在blockingGet()这里等待onComplete,根本没法继续执行doOnNext的后续代码,更没法让源发送onComplete。这就形成了死锁:main线程等onComplete,onComplete等main线程完成doOnNext,所以一直阻塞。
singleElement()和singleOrError()逻辑一致,只是后者在元素数量不符合时抛出错误,两者都需要等待onComplete来验证元素数量,因此都会陷入同样的死锁。
2. firstElement()为什么不会阻塞?
firstElement()的设计逻辑完全不同:它只关心第一个发射的元素,只要收到第一个元素,就会立即把元素发射出去,并且自己进入完成状态,完全不需要等待源的onComplete信号。
所以当你调用cachedFlowable.firstElement().blockingGet()时:
- 它会立即拿到缓存的
1,直接返回,不会等待后续的onComplete。这样doOnNext就可以继续执行,打印after blockingGet 1,之后第一个订阅的流程继续推进,源发送onComplete,整个流程正常结束。
内容的提问来源于stack exchange,提问作者Courageous Dilemma
相关产品推荐
相关产品推荐

