RxJava中SwitchMap结合延迟Observable无输出问题咨询
问题分析与解决方案
你的代码有两个核心问题导致没有输出,咱们一步步拆解清楚:
1. switchMap逻辑与无限interval的冲突
你对switchMap的核心理解是对的——它会取消之前订阅的Observable,只保留最新的那个,但问题出在Observable.interval(2, TimeUnit.SECONDS)是无限发射事件的:
- 第0秒发射
0,switchMap创建一个需要延迟5秒才输出0的Observable; - 第2秒发射
1,switchMap立即取消那个还没到延迟时间的Observable,转而订阅新的、输出2的Observable; - 第4秒发射
2,再次取消上一个Observable; - 以此类推,永远不会有Observable能走完5秒的延迟——因为每2秒就会有新的事件打断它,自然不会有任何输出。
你预期20秒后输出8,应该是默认假设interval会在某个点停止发射,但默认的interval是无限序列,除非手动终止它。
2. 主线程退出导致守护线程被终止
另外,如果你是在普通main方法中运行这段代码,主线程执行完subscribe后就会直接结束,而RxJava的interval默认运行在Schedulers.computation()线程池,这些线程是守护线程——当JVM中所有非守护线程结束时,守护线程会被强制终止,所以即使有机会输出,程序也会提前退出。
修正方案
如果你想看到预期的8输出,可以做以下调整:
方案1:让interval只发射有限次数
比如让interval只发射到第4个事件(对应输出8),然后终止:
Observable.interval(2, TimeUnit.SECONDS) .take(5) // 只取前5个事件:0,1,2,3,4 .switchMap(integer -> Observable .just(integer * 2) .delay(5, TimeUnit.SECONDS)) .subscribe(integer -> { System.out.println(integer);}); // 阻塞主线程,确保程序不提前退出(main方法中需要) try { Thread.sleep(15000); // 等待足够时间让最后一个delay完成 } catch (InterruptedException e) { e.printStackTrace(); }
这里take(5)会让interval发射0-4后终止,最后一个事件4对应的Observable会在第8+5=13秒时输出8,主线程sleep15秒能确保程序等到输出完成。
方案2:用concatMap替代switchMap(如果需要保留所有事件)
如果你不想取消之前的Observable,而是希望每个事件都能输出,那应该用concatMap——它会按顺序执行每个Observable,不会取消之前的:
Observable.interval(2, TimeUnit.SECONDS) .take(5) .concatMap(integer -> Observable .just(integer * 2) .delay(5, TimeUnit.SECONDS)) .subscribe(integer -> { System.out.println(integer);}); try { Thread.sleep(30000); // 5个事件,每个延迟5秒,总耗时约27秒 } catch (InterruptedException e) { e.printStackTrace(); }
这样会依次输出0,2,4,6,8,每个间隔5秒。
内容的提问来源于stack exchange,提问作者Mykola Chykin
相关产品推荐
相关产品推荐

