Project Reactor 3中publishOn与subscribeOn的差异、优先级及使用场景咨询
关于Flux中publishOn与subscribeOn的疑问及解答
首先得解决你遇到的日志无输出问题——这跟两个操作符的优先级没关系,纯粹是因为你用了异步调度器,但主线程没等异步任务执行完就直接结束了。你可以在代码最后加个blockLast()(比Thread.sleep()更优雅)来等待Flux执行完成,这样就能看到日志输出了。
接下来咱们聊聊这两个操作符的核心区别、适用场景,以及你提到的优先级问题:
一、核心区别:管的是不同的执行阶段
- subscribeOn:它负责整个Flux链的上游生产环节——不管你把它放在链的哪个位置,它只会决定「订阅触发后,上游的数据生成、map这类前置操作」在哪个线程池运行。比如你例子里的
subscribeOn(Schedulers.parallel()),是让Flux.just(1,2,3,4)和map(i->i*2)都在parallel线程池执行,但因为主线程没等,这些异步任务还没来得及输出日志就被终止了。 - publishOn:它负责它之后的下游消费环节——只要是在它后面的操作符、订阅逻辑,都会切换到指定的线程池执行。你单独用
publishOn时,上游的操作(Flux.just和map)其实是在主线程同步执行的,主线程会等上游做完,所以能看到完整日志;而下游的subscribe(elements::add)则是在elastic线程池运行的。
二、优先级与执行逻辑:不存在谁“更推荐”,各司其职
这俩操作符没有谁优先级更高的说法,而是各自负责不同的阶段:
- 如果同时使用,
subscribeOn管上游生产的线程,publishOn管它之后的下游线程。比如你的例子里,map在parallel线程执行,subscribe在elastic线程执行,但主线程提前结束,导致异步任务没跑完。 - 额外提个细节:
subscribeOn在整个链里只会生效一次,哪怕你加多个,只有第一个起作用;而publishOn可以加多次,每次添加都会切换之后的执行线程。
三、适用场景:按需选择
- 用subscribeOn的场景:当你上游有阻塞性操作(比如数据库查询、文件读取、调用第三方接口)时,用它把这些操作放到异步线程池(比如
Schedulers.boundedElastic()),避免阻塞主线程。比如从数据库拉取数据的Flux,就适合用subscribeOn把查询逻辑放到异步线程。 - 用publishOn的场景:当你下游有耗时的消费/处理操作(比如复杂计算、写入文件、UI更新)时,用它切换到合适的线程池;或者需要在不同阶段切换线程适配操作类型——比如上游在IO线程拿数据,用
publishOn(Schedulers.parallel())切换到计算线程做数据处理,再用publishOn(Schedulers.boundedElastic())切回IO线程写文件。
最后给你修改后的代码示例,保证能看到日志:
System.out.println("*********Calling Concurrency************"); List<Integer> elements = new ArrayList<>(); Flux.just(1, 2, 3, 4) .map(i -> i * 2) .log() .publishOn(Schedulers.elastic()) .subscribeOn(Schedulers.parallel()) .doOnNext(elements::add) .blockLast(); // 等待Flux执行完成,避免主线程提前退出 System.out.println("-------------------------------------");
内容的提问来源于stack exchange,提问作者KayV
相关产品推荐
相关产品推荐

