为何我的RxJava订阅回调未执行?异步订阅问题求助
RxJava中耗时Observable未触发onNext,主线程提前结束的问题
我来帮你分析下这个问题哈~你遇到的核心问题是后台线程与主线程的生命周期不同步导致的:
你把耗时计算放在fromCallable里,通过subscribeOn(Schedulers.io())把任务切换到IO后台线程执行,但调用test()的主线程在触发订阅后就直接走完方法流程了。如果这是普通的Java/Kotlin控制台程序,主线程结束后整个程序就会退出,IO线程还没来得及完成计算并触发onNext回调,自然就看不到输出"6"了。
而如果不用IO调度器,默认会在当前线程(也就是主线程)执行fromCallable里的逻辑,Thread.sleep(4000)会直接阻塞主线程,直到计算完成才会继续执行onNext回调,所以能正常输出结果。
下面给你几个不同场景下的可行解决方案:
方案1:阻塞式订阅(测试场景专用)
把普通的subscribe()换成blockingSubscribe(),它会让当前线程(这里是主线程)阻塞,直到Observable的整个生命周期完成(包括onNext、onComplete):
fun test(){ fetchNumber(2,4) .subscribeOn(Schedulers.io()) .doOnSubscribe { println("subscribed") } .blockingSubscribe({ println(it)}) } private fun fetchNumber(a: Int, b: Int) : Observable<Int> { return Observable.fromCallable { Thread.sleep(4000) a + b } }
方案2:让主线程保持存活(快速测试用)
在test()方法最后加一段休眠,给IO线程足够的时间完成计算:
fun test(){ fetchNumber(2,4) .subscribeOn(Schedulers.io()) .doOnSubscribe { println("subscribed") } .subscribe({ println(it)}) // 休眠时间要比耗时任务的时间长一点 Thread.sleep(5000) }
方案3:用CountDownLatch优雅等待(推荐测试场景使用)
这是更专业的测试方式,不会像固定休眠那样浪费时间,任务完成后会立即唤醒主线程:
import java.util.concurrent.CountDownLatch fun test(){ val latch = CountDownLatch(1) fetchNumber(2,4) .subscribeOn(Schedulers.io()) .doOnSubscribe { println("subscribed") } .doOnComplete { latch.countDown() } // 任务完成时触发计数减一 .subscribe({ println(it)}) latch.await() // 主线程等待,直到计数归0 } private fun fetchNumber(a: Int, b: Int) : Observable<Int> { return Observable.fromCallable { Thread.sleep(4000) a + b } }
补充:Android等有Looper的环境无需处理
如果是在Android应用中使用这段代码,主线程本身是带有Looper循环的,不会执行完就退出,IO线程完成计算后会正常触发onNext回调,所以不会出现这个问题。
内容的提问来源于stack exchange,提问作者dev
相关产品推荐
相关产品推荐

