You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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退出规则、不同运行环境的执行逻辑共同导致的:

  1. JVM退出的基础规则
    只有当JVM中所有存活的非守护线程都执行结束时,进程才会正常退出;守护线程的存活不会阻止JVM退出。
  2. 两种parallel调度器的核心差异
    • 调用Schedulers.newParallel(xxx)时,会创建一个全新的独立并行调度器实例,它底层线程池生成的工作线程默认是非守护线程,且默认配置下核心线程不会因空闲被回收。当Flux流的3个元素全部消费完成后,代码中创建的sub、pub两个调度器的工作线程依然处于存活、等待新任务的状态,属于存活的非守护线程,直接导致JVM判定进程不能退出,表现为程序持续挂起。
    • 调用Schedulers.parallel()时,拿到的是Reactor全局共享的并行调度器实例,它底层的工作线程全部是守护线程,哪怕这些线程在流结束后依然存活,也不会阻止JVM退出,因此主线程跑完后进程可以正常终止。
  3. 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.01 04:24:23