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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 07:12:50