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

Project Reactor热发布者:首订阅者无元素、次订阅者重复消费原因

Project Reactor热发布者异常现象成因分析

演示代码

@Test
void testSimpleFlux() {
    Flux<Integer> intFlux = Flux.just(1, 2, 3, 4, 5)
            .publish()
            .autoConnect()
            .log();
    intFlux.subscribeOn(Schedulers.newParallel("Thread-one"))
            .subscribe(System.out::println);
    intFlux.subscribeOn(Schedulers.newParallel("Thread-two"))
            .subscribe(System.out::println);
}

预期结果

首个订阅者因先订阅并触发元素发射,会消费全部元素;第二个订阅者因订阅较晚,仅能消费部分元素(例如2、3、4、5)。

实际结果

首个订阅者未接收到任何元素,第二个订阅者则重复消费了每个整数两次,日志如下:

08:10:04.095 [Thread-one-1] INFO reactor.Flux.AutoConnect.1 -- onSubscribe(FluxPublish.PublishInner)
08:10:04.095 [Thread-two-2] INFO reactor.Flux.AutoConnect.1 -- onSubscribe(FluxPublish.PublishInner)
08:10:04.106 [Thread-two-2] INFO reactor.Flux.AutoConnect.1 -- request(unbounded)
08:10:04.106 [Thread-one-1] INFO reactor.Flux.AutoConnect.1 -- request(unbounded)
08:10:04.108 [Thread-two-2] INFO reactor.Flux.AutoConnect.1 -- onNext(1)
1
08:10:04.109 [Thread-two-2] INFO reactor.Flux.AutoConnect.1 -- onNext(1)
1
08:10:04.109 [Thread-two-2] INFO reactor.Flux.AutoConnect.1 -- onNext(2)
2
08:10:04.109 [Thread-two-2] INFO reactor.Flux.AutoConnect.1 -- onNext(2)
2
08:10:04.109 [Thread-two-2] INFO reactor.Flux.AutoConnect.1 -- onNext(3)
3
08:10:04.109 [Thread-two-2] INFO reactor.Flux.AutoConnect.1 -- onNext(3)
3
08:10:04.110 [Thread-two-2] INFO reactor.Flux.AutoConnect.1 -- onNext(4)
4
08:10:04.110 [Thread-two-2] INFO reactor.Flux.AutoConnect.1 -- onNext(4)
4
08:10:04.110 [Thread-two-2] INFO reactor.Flux.AutoConnect.1 -- onNext(5)
5
08:10:04.111 [Thread-two-2] INFO reactor.Flux.AutoConnect.1 -- onNext(5)
5
08:10:04.111 [Thread-two-2] INFO reactor.Flux.AutoConnect.1 -- onComplete()
08:10:04.112 [Thread-two-2] INFO reactor.Flux.AutoConnect.1 -- onComplete()

异常成因

  1. subscribeOn位置错误导致线程复用:你把subscribeOn放在了publish().autoConnect()之后,这只会让订阅热流的操作在指定后台线程执行,而Flux.just是同步源,它的元素发射线程由第一个完成订阅流程的线程决定(日志中是Thread-two)。由于没有用publishOn指定消费线程,两个订阅者的消费逻辑会直接复用源的发射线程,导致所有输出都显示在Thread-two上,让你误以为只有第二个订阅者收到了元素。

  2. 订阅时机同步导致无元素遗漏:主线程会快速连续执行两个subscribe调用,subscribeOn指定的两个后台线程几乎同时完成了订阅和request(unbounded)操作——在Flux.just开始发射元素前,两个订阅者都已完成订阅流程并准备就绪。因此autoConnect触发源发射后,热流会将所有元素推送给两个订阅者,每个元素被输出两次。

  3. 对autoConnect触发逻辑的误解:你预期第一个订阅者会先触发源发射,导致第二个订阅者错过早期元素,但实际上subscribeOn将订阅操作异步到后台线程,主线程不会等待第一个订阅完成就执行第二个订阅,两个后台线程的订阅流程几乎同步完成,源还未开始发射就已有两个订阅者等待,因此没有元素被遗漏。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 21:34:52