RxJava中concatMap内Observable加delay后无数据发射的原因咨询
为什么concatMap内添加delay后没有数据输出?
这问题其实挺典型的,我来给你拆解背后的运行逻辑:
先看两段代码的核心差异
Snippet1(无delay)
Observable.range(1, 10) .concatMap { Observable.just(it) } .subscribe{ print(it)}
这段代码里所有操作都是同步执行的:
Observable.range在当前线程(比如主线程)直接发射1到10的数字concatMap把每个数字转成Observable.just(it),这个just也是同步发射数据subscribe的回调同样在当前线程执行,所有数字会被连续打印,主线程要等所有操作完成才会继续往下走,自然能看到完整输出。
Snippet2(带delay)
Observable.range(1, 10) .concatMap { Observable.just(it).delay(1, TimeUnit.SECONDS) } .subscribe{ print(it)}
这里的关键问题出在**delay的线程调度和主线程的生命周期**上:
delay默认会把数据发射任务切换到Schedulers.computation()线程池执行,也就是异步延迟1秒后再发射数据- 你的主线程(比如是main函数或简单测试代码)在调用
subscribe后,没有任何阻塞等待逻辑,会直接执行完所有同步代码然后退出 - 当主线程退出时,整个JVM进程会终止,所有后台线程(包括执行delay任务的computation线程)都会被强制停止,根本等不到1秒后的数据发射,所以看不到任何输出。
另外要提concatMap的特性:它是串行处理每个内部Observable的,必须等前一个内部Observable完成,才会处理下一个。但在这个场景里,第一个内部Observable还没来得及发射数据,主线程就已经退出了,后面的流程根本没机会执行。
怎么解决这个问题?
要让延迟后的Observable正常发射数据,核心是让主线程不要提前退出,等到整个流完成再结束。举个实际例子:
import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; import io.reactivex.Observable; public class RxTest { public static void main(String[] args) throws InterruptedException { CountDownLatch latch = new CountDownLatch(1); Observable.range(1, 10) .concatMap(it -> Observable.just(it).delay(1, TimeUnit.SECONDS)) .doOnComplete(latch::countDown) // 流完成时触发latch计数减一 .subscribe(System.out::print); latch.await(); // 主线程阻塞,等待流完成 } }
这样主线程会一直等待,直到所有10个数据都发射完成,就能看到依次输出1到10的结果了。
如果是在Android这类有主线程Looper循环的环境里,主线程本身会一直运行(除非应用退出),这时候加delay是能正常看到输出的,因为后台线程的延迟任务可以正常执行并回调到主线程。
内容的提问来源于stack exchange,提问作者Suraj
相关产品推荐
相关产品推荐

