Java Reactor 3中编写的Flux.create相关代码无输出结果该如何解决?
Java Reactor 3 代码无输出问题分析与解决方案
问题原因
- 核心原因:
publishOn(Schedulers.elastic())会将后续的消费逻辑调度到守护线程池中执行。你的主线程在调用subscribe()方法之后就直接运行结束了,JVM检测到当前没有存活的非守护线程时会直接终止进程,所有还没执行的守护线程任务都会被丢弃,所以看不到打印输出。 - 次要问题:你注释了
sink.complete()方法,会导致当前Flux一直处于未完成状态,即使主线程不退出,也永远不会触发subscribe的第三个完成回调。
解决办法
你可以根据实际使用场景选择以下任意一种方案:
方案1:阻塞等待流处理完成(适合简单调试、同步场景)
先打开注释的sink.complete()给流添加明确的结束信号,再调用blockLast()方法让主线程阻塞直到流处理结束:
Flux.create(sink -> { sink.next("produce a number: " + Math.random() * 100); sink.complete(); }).publishOn(Schedulers.elastic()) .subscribe( consumer -> System.out.println(Thread.currentThread().getName() + consumer), error -> System.out.println("error!" + error), () -> System.out.println("task complete!") ) .blockLast();
方案2:手动控制主线程等待(适合流不需要立刻结束的场景)
用CountDownLatch手动控制主线程的退出时机,等到你需要的业务逻辑执行完再放行主线程:
// 初始化计数器,这里设置为1代表只需要等待1个执行完成信号 CountDownLatch latch = new CountDownLatch(1); Flux.create(sink -> { sink.next("produce a number: " + Math.random() * 100); // 如果你后续还要生产更多元素,可以先不调用sink.complete() }).publishOn(Schedulers.elastic()) .subscribe( consumer -> { System.out.println(Thread.currentThread().getName() + consumer); latch.countDown(); // 元素处理完成后计数器减1 }, error -> { System.out.println("error!" + error); latch.countDown(); }, () -> { System.out.println("task complete!"); latch.countDown(); } ); // 主线程阻塞,直到计数器变为0再继续执行 latch.await();
方案3:单元测试场景用StepVerifier
如果是编写单元测试验证流行为,可以用Reactor官方提供的StepVerifier工具,它会自动阻塞等待流处理完成,不需要手动控制线程:
Flux<String> testFlux = Flux.create(sink -> { sink.next("produce a number: " + Math.random() * 100); sink.complete(); }).publishOn(Schedulers.elastic()); StepVerifier.create(testFlux) .expectNextCount(1) // 校验是否收到1个元素 .expectComplete() // 校验是否收到完成信号 .verify(); // 自动阻塞直到流处理结束
内容的提问来源于stack exchange,提问作者光明的心
相关产品推荐
相关产品推荐

