在Project Reactor中如何将Mono多播给多个订阅者
基于Mono实现热发布器的两种常用方案
方案1:使用publish()手动控制触发时机(和Flux多播逻辑完全对齐)
你可以直接调用Mono的publish()方法得到MonoProcessor,和你用到的ConnectableFlux逻辑一致,支持提前注册多个订阅者,手动触发上游执行,上游逻辑仅运行一次,结果会广播给所有已注册的订阅者。
示例代码如下:
// 原始冷Mono,替换为你的业务逻辑即可 Mono<String> coldMono = Mono.fromSupplier(() -> { System.out.println("上游Mono逻辑仅执行一次"); return "业务查询结果"; }).log(); // 转换为可连接的热发布器 MonoProcessor<String> hotMono = coldMono.publish(); // 注册多个订阅者 hotMono.map(res -> "订阅者1处理结果:" + res) .subscribe(System.out::println); hotMono.map(res -> "订阅者2处理结果:" + res) .subscribe(System.out::println); // 手动触发上游执行,结果多播给所有已注册的订阅者 hotMono.connect();
方案2:使用cache()自动缓存结果(更适合大部分业务场景)
如果你不需要精确控制上游的触发时机,只需要保证上游Mono仅执行一次,后续所有订阅者都能拿到复用结果,可以直接使用cache()方法,实现更简洁:
// 转换为带缓存的热发布器,上游执行后结果会永久缓存 Mono<String> cachedHotMono = Mono.fromSupplier(() -> { System.out.println("上游Mono逻辑仅执行一次"); return "业务查询结果"; }).log() .cache(); // 第一次订阅会触发上游执行 cachedHotMono.map(res -> "订阅者1处理结果:" + res) .subscribe(System.out::println); // 后续所有订阅直接读取缓存结果,不会重新触发上游逻辑 cachedHotMono.map(res -> "订阅者2处理结果:" + res) .subscribe(System.out::println);
如果需要设置缓存有效期,可以传入时间参数,例如cache(Duration.ofMinutes(5))表示结果缓存5分钟后失效,失效后新的订阅会重新触发上游执行。
内容的提问来源于stack exchange,提问作者n0noob
相关产品推荐
相关产品推荐

