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
相关产品推荐
相关产品推荐

