如何在Reactor Java中实现单源多生产者的设计?
Project Reactor 单数据源多生产者实现方案
你现在的代码有两个问题:一是用map同步调用生产者的send方法,会阻塞当前线程,不符合Reactor异步非阻塞的设计思路;二是无法处理send可能返回的异步结果,也没法保证两个发送操作的完成性。
下面给几种符合Reactor规范的实现方式:
方式一:确保两个发送操作都完成(flatMap+Mono.when)
如果firstProducer.send和secondProducer.send返回Mono<Void>(代表异步发送操作),可以用Mono.when合并两个异步任务,等待所有操作完成后再继续流程:
private Mono<String> main() { return Mono.just("simple message") .flatMap(msg -> Mono.when( firstProducer.send(msg), secondProducer.send(msg) ).thenReturn(msg) ); }
Mono.when会等待所有传入的Mono执行完成后,再发出完成信号thenReturn(msg)用来在两个发送操作完成后返回原消息
方式二:发后即忘(doOnNext)
如果不需要等待发送操作完成,只是触发发送即可,用doOnNext不会阻塞主流程:
private Mono<String> main() { return Mono.just("simple message") .doOnNext(msg -> firstProducer.send(msg).subscribe()) .doOnNext(msg -> secondProducer.send(msg).subscribe()); }
注意:这里必须调用
subscribe(),否则doOnNext里的Mono不会实际执行。这种方式适合无需关心发送结果的场景,发送失败也不会影响主流程。
方式三:广播模式(publish+subscribe)
如果需要将数据源事件广播给多个生产者,用publish创建可连接的流,再分别订阅给不同生产者:
private Mono<String> main() { Mono<String> source = Mono.just("simple message"); // 给第一个生产者订阅数据源 source.subscribe(msg -> firstProducer.send(msg).subscribe()); // 给第二个生产者订阅数据源 source.subscribe(msg -> secondProducer.send(msg).subscribe()); return source; }
如果是冷数据源(比如Mono.just),每个订阅会重新触发一次数据源逻辑;如果想让多个订阅者共享同一份事件,用publish().autoConnect():
private Mono<String> main() { ConnectableFlux<String> source = Mono.just("simple message").publish(); source.subscribe(msg -> firstProducer.send(msg).subscribe()); source.subscribe(msg -> secondProducer.send(msg).subscribe()); source.connect(); // 手动触发数据源发射事件 return source.next(); }
原代码改进建议
如果你的send是同步方法(不返回Mono/Flux),至少要用doOnNext代替map——map的语义是转换数据,doOnNext才是用来执行副作用操作的:
private Mono<String> main() { return Mono.just("simple message") .doOnNext(firstProducer::send) .doOnNext(secondProducer::send); }
内容的提问来源于stack exchange,提问作者kaimada
相关产品推荐
相关产品推荐

