Project Reactor中Schedulers.newParallel()执行完Flux后为何挂起不终止?
问题描述
现有一个元素类型为String的基础Flux,在main()方法中运行如下代码:
package com.example; import reactor.core.publisher.Flux; import reactor.core.scheduler.Schedulers; import reactor.util.Logger; import reactor.util.Loggers; import java.util.Arrays; import java.util.List; public class Parallel { private static final Logger log = Loggers.getLogger(Parallel.class.getName()); private static List<String> COLORS = Arrays.asList("red", "white", "blue"); public static void main(String[] args) throws InterruptedException { Flux<String> flux = Flux.fromIterable(COLORS); flux .log() .map(String::toUpperCase) .subscribeOn(Schedulers.newParallel("sub")) .publishOn(Schedulers.newParallel("pub", 1)) .subscribe(value -> { log.info("==============Consumed: " + value); }); } }
运行时出现如下现象:
- 程序会持续挂起,无法自动停止,必须手动终止进程
- 如果将代码中的
.newParallel()替换为.parallel(),程序就能按预期正常执行、正常退出 - 如果将这段代码作为JUnit测试用例运行,则不会出现挂起问题,可正常执行完成
核心疑问:为什么程序无法在Flux元素全部发射完成后自行结束运行?出现该挂起现象的根本原因是什么?
根本原因分析
这个挂起问题和Flux流本身的元素发射、消费逻辑无关,是Reactor调度器的线程属性、JVM退出规则、不同运行环境的执行逻辑共同导致的:
- JVM退出的基础规则
只有当JVM中所有存活的非守护线程都执行结束时,进程才会正常退出;守护线程的存活不会阻止JVM退出。 - 两种parallel调度器的核心差异
- 调用
Schedulers.newParallel(xxx)时,会创建一个全新的独立并行调度器实例,它底层线程池生成的工作线程默认是非守护线程,且默认配置下核心线程不会因空闲被回收。当Flux流的3个元素全部消费完成后,代码中创建的sub、pub两个调度器的工作线程依然处于存活、等待新任务的状态,属于存活的非守护线程,直接导致JVM判定进程不能退出,表现为程序持续挂起。 - 调用
Schedulers.parallel()时,拿到的是Reactor全局共享的并行调度器实例,它底层的工作线程全部是守护线程,哪怕这些线程在流结束后依然存活,也不会阻止JVM退出,因此主线程跑完后进程可以正常终止。
- 调用
- JUnit场景无挂起的原因
JUnit的测试执行器不会等待异步派生的非守护线程运行结束,只要测试方法本身的同步代码执行完毕,测试执行器就会直接退出测试进程,因此感知不到挂起问题。
修复方案
如果要在main方法中使用newParallel()又不想进程挂起,只需要在流完全结束后主动释放调度器资源即可,示例代码如下:
public static void main(String[] args) { Flux<String> flux = Flux.fromIterable(COLORS); // 提前声明自定义调度器实例 Scheduler subScheduler = Schedulers.newParallel("sub"); Scheduler pubScheduler = Schedulers.newParallel("pub", 1); flux .log() .map(String::toUpperCase) .subscribeOn(subScheduler) .publishOn(pubScheduler) .doFinally(signal -> { // 流终止(正常完成、异常、取消)时主动销毁调度器,关闭线程池 subScheduler.dispose(); pubScheduler.dispose(); }) .subscribe(value -> { log.info("==============Consumed: " + value); }); }
流执行完成后两个自定义调度器的线程池会被主动关闭,所有非守护线程终止,JVM就可以正常退出了。
内容的提问来源于stack exchange,提问作者Roman Y
相关产品推荐
相关产品推荐

