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()
异常成因
subscribeOn位置错误导致线程复用:你把subscribeOn放在了publish().autoConnect()之后,这只会让订阅热流的操作在指定后台线程执行,而Flux.just是同步源,它的元素发射线程由第一个完成订阅流程的线程决定(日志中是Thread-two)。由于没有用publishOn指定消费线程,两个订阅者的消费逻辑会直接复用源的发射线程,导致所有输出都显示在Thread-two上,让你误以为只有第二个订阅者收到了元素。订阅时机同步导致无元素遗漏:主线程会快速连续执行两个
subscribe调用,subscribeOn指定的两个后台线程几乎同时完成了订阅和request(unbounded)操作——在Flux.just开始发射元素前,两个订阅者都已完成订阅流程并准备就绪。因此autoConnect触发源发射后,热流会将所有元素推送给两个订阅者,每个元素被输出两次。对
autoConnect触发逻辑的误解:你预期第一个订阅者会先触发源发射,导致第二个订阅者错过早期元素,但实际上subscribeOn将订阅操作异步到后台线程,主线程不会等待第一个订阅完成就执行第二个订阅,两个后台线程的订阅流程几乎同步完成,源还未开始发射就已有两个订阅者等待,因此没有元素被遗漏。
内容的提问来源于stack exchange,提问作者Powet
相关产品推荐
相关产品推荐

