RxJava中Schedulers未将任务切换至IO线程的问题排查
看你的日志结果,Inside Map和Inside subscribe始终跑在NIO线程里,完全没切换到期望的RxCachedSchedulerThread,核心问题是混淆了RxJava中subscribeOn和observeOn的作用,再加上上游异步操作的线程绑定导致subscribeOn没生效。
为什么subscribeOn(Schedulers.io())没起作用?
subscribeOn的核心作用是指定Observable链最上游的订阅执行线程,但你的products Single来自recommendationService.getRecommendations(ids),里面的flatMap调用了productClient.getProduct(id)——这个HTTP客户端操作(看起来是Micronaut的客户端)默认是在NIO EventLoop线程执行异步回调的。也就是说,整个Single的数据流发射已经被绑定到了NIO线程,下游的subscribeOn根本无法覆盖上游已经确定的发射线程。
另外要注意:subscribeOn只对订阅动作发生时的初始线程生效,如果上游已经有异步操作切换了线程,它就无法影响后续的下游操作了。
解决方法:用observeOn(Schedulers.io())切换下游线程
如果你希望map和subscribe里的逻辑跑在IO线程,应该使用observeOn——它的作用就是强制下游所有操作(包括map、subscribe)切换到指定线程执行,不管上游在哪个线程。
修改你的代码如下:
log.info("NIO Thread : " + Thread.currentThread().getName()); Single<List<Product>> products = recommendationService.getRecommendations(ids); log.info("NIO Response Thread : " + Thread.currentThread().getName()) ; // 先subscribeOn(可选,若要让上游订阅也跑在IO线程),再用observeOn切换下游到IO线程 products.subscribeOn(Schedulers.io()) .observeOn(Schedulers.io()) // 关键:添加这一行切换下游线程 .map(s -> { log.info("Inside Map : " + Thread.currentThread().getName()); return s; }) .subscribe( s -> log.info("Inside subscribe :" + s.get(0).getName())); log.info("After thread "); return products;
额外优化:让getRecommendations里的操作也跑在IO线程
如果连RecommendationService.getRecommendations()里的flatMap逻辑也想放到IO线程,应该在getRecommendations的Observable链中指定subscribeOn:
public Single<List<Product>> getRecommendations(List<String> ids) throws MalformedURLException { log.info("RecommendationService.getRecommendations() called..."); return Observable.fromIterable(ids) .subscribeOn(Schedulers.io()) // 让整个Observable链的源头跑在IO线程 .flatMap(id -> productClient.getProduct(id).toObservable()) .toList(); }
这样一来,结合下游的observeOn,整个数据流的处理就都能在IO线程执行了。
验证效果
修改后你会看到日志变成:
16:17:52.097 [RxCachedSchedulerThread-xx] INFO com.igaurav.RecommendationController - Inside Map : RxCachedSchedulerThread-xx 16:17:52.097 [RxCachedSchedulerThread-xx] INFO com.igaurav.RecommendationController - Inside subscribe :Micronaut in Action
内容的提问来源于stack exchange,提问作者Ishant Gaurav

