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

Reactor Core:多订阅者场景下无限Hot Flux运行异常问题

问题原因与解决逻辑拆解

嘿,这个问题我之前玩Reactor的时候也踩过坑,咱们一步步捋清楚背后的逻辑:

1. Reactor线程的「守护属性」是核心根源

Reactor框架默认使用的调度器(比如Schedulers.parallel()、Schedulers.single())创建的线程,全都是守护线程。而Java里守护线程的特性很关键:当JVM中所有非守护线程都终止时,JVM会直接退出,完全不会等待守护线程完成任务。

2. 你的测试代码的线程运行逻辑

咱们还原下你的测试场景:

  • 你的测试代码(不管是main方法还是JUnit测试方法)运行在非守护线程里(比如主线程、JUnit的测试线程)。
  • 当你创建「hot无限Flux」(比如用Flux.interval()搭配publish().autoConnect())并订阅后,Flux的事件生成、数据推送逻辑,其实是跑在Reactor的守护线程上的。
  • 如果测试代码没有任何阻塞逻辑,主线程会快速执行完所有代码然后终止。这时JVM发现没有剩余的非守护线程了,就会直接终止整个进程——那些在守护线程上跑的Flux流自然被强制停掉,这就是为什么第二个订阅者刚订阅,进程就直接停止了。

3. Thread.currentThread().join()的解决原理

你加的这句代码本质是让主线程(非守护线程)一直处于阻塞状态,不让它终止:

  • Thread.currentThread()获取的就是当前执行这段代码的线程——也就是你的测试主线程(非守护线程)。
  • 调用join()方法的作用是「让调用线程等待目标线程终止」,这里调用线程和目标线程都是主线程,相当于让主线程无限等待自己终止(这显然不可能发生),所以主线程会一直阻塞在这里。
  • 只要主线程不终止,JVM就不会退出,Reactor的守护线程就能持续运行,两个订阅者自然就能一起接收流数据了。

小补充:更规范的测试姿势

如果是写单元测试,其实不推荐用join()这种手动阻塞的方式,Reactor官方提供的StepVerifier更适合测试无限流,比如:

StepVerifier.create(yourHotFlux)
    .expectNextCount(3) // 验证接收3条数据
    .thenCancel() // 主动取消订阅,避免资源泄漏
    .verify();

这种方式可以精确控制测试流程,不需要手动阻塞线程,测试完成后会自动清理资源。

内容的提问来源于stack exchange,提问作者Cepr0

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 09:04:21