RxJava中observeOn位置不同为何行为异常?无背压且单线程?
RxJava线程复用与背压表现疑惑解答
先梳理你的两个场景和核心疑问:
场景1:observeOn置于map之前的表现
你的代码结构如下:
Disposable x = pro .generateFlowable( new File("path\\to\\file.raw")) .subscribeOn(Schedulers.io(), false) .observeOn(Schedulers.io()) .map(y -> { System.out.println(Thread.currentThread().getName() + " xxx"); return y; }) .subscribe(onNext -> { System.out.println(Thread.currentThread().getName() + " " + new String(onNext)); Thread.sleep(100); }, Throwable::printStackTrace, () -> { System.out.println("Done"); t.end(); System.out.println(t.getTotalTime()); });
运行时输出交替出现RxCachedThreadScheduler-1 xxx和RxCachedThreadScheduler-1 Line1,全程复用同一个线程。
场景2:observeOn移至subscribe之前的表现
调整后的代码结构:
Disposable x = pro .generateFlowable( new File("path\\to\\file.raw")) .subscribeOn(Schedulers.io(), false) .map(y -> { System.out.println(Thread.currentThread().getName() + " xxx"); return y; }) .observeOn(Schedulers.io()) .subscribe(onNext -> { System.out.println(Thread.currentThread().getName() + " " + new String(onNext)); Thread.sleep(100); }, Throwable::printStackTrace, () -> { System.out.println("Done"); t.end(); System.out.println(t.getTotalTime()); });
此时输出变为先批量出现RxCachedThreadScheduler-1 xxx,再批量出现RxCachedThreadScheduler-1 Line1,线程仍为同一个。
你的核心疑问:
- 为何会出现交替/批量输出的差异?
- 为何全程只复用一个线程?
observeOn看起来没切换线程,它到底生效了吗?
1. 线程复用的本质:Schedulers.io()的线程池机制
Schedulers.io()背后是一个可缓存的线程池,它的核心逻辑是优先复用空闲线程,只有当所有线程都处于忙碌状态时才会创建新线程。
在你的两个场景中:
- 场景1里,
subscribeOn(Schedulers.io())指定上游的generate和后续的map在io线程池执行,而observeOn(Schedulers.io())又指定下游的subscribe回调也在同一个io线程池。当下游处理(Thread.sleep(100))结束释放线程后,上游的map任务刚好可以复用这个空闲线程,所以你看到的是同一个线程ID。 - 场景2里,
subscribeOn负责上游generate和map的线程,observeOn负责下游subscribe的线程,但因为是同一个线程池,当下游线程空闲时,上游的map任务依然可以复用它,所以线程ID还是同一个。
如果想验证线程切换效果,你可以把observeOn的调度器换成Schedulers.newThread()(每次都会创建新线程),这时就能看到map和subscribe使用不同的线程了。
2. 交替/批量输出的背后:背压机制的作用
你的generateFlowable用了Flowable.generate,它是原生支持背压的操作符,会严格根据下游的请求量来生产数据。
- 场景1的交替输出:
observeOn在map之前,意味着map和subscribe都在同一个线程执行。当下游subscribe处理完一条数据(睡眠100ms)后,会向上游请求下一条数据,上游的map立刻执行并发送数据,所以两者交替输出。 - 场景2的批量输出:
observeOn在map之后,此时observeOn会创建一个默认大小为128的缓冲区。当下游subscribe处理速度很慢时,上游的map会先生产填满缓冲区的数据(批量输出xxx),然后暂停生产;等下游慢慢处理完缓冲区里的数据,再向上游请求新数据,这时又会批量处理一批,所以你看到先批量xxx再批量Line1——这正是背压生效的典型表现。
3. observeOn的生效逻辑:你误解了它的作用
observeOn的核心作用是指定其下游所有操作的执行线程,和上游操作符无关:
- 场景1中,
observeOn在map之前,所以map和subscribe都在observeOn指定的线程执行; - 场景2中,
observeOn在map之后,所以只有subscribe在observeOn的线程执行,map则在subscribeOn的线程执行。
之所以看起来没切换线程,只是因为两个调度器用了同一个线程池,线程被复用了而已,observeOn其实已经正常生效了。
内容的提问来源于stack exchange,提问作者chhil
相关产品推荐
相关产品推荐

