Reactor应用中如何等待两个发布者完成并返回其一结果
解决方案
首先要解决一个关键问题:你的代码里没有在Mono的发送逻辑结束后调用dataSink.tryEmitComplete(),这会导致Flux永远不会进入完成状态,后续的等待逻辑永远无法结束。所以第一步需要修正Mono的实现:
var mono = Mono.fromSupplier(() -> { HashSet<Long> map = new HashSet<>(); for (int i = 0; i < data.length; i++) { Long datum = data[i]; map.add(datum); log.info("sending....{}", i); dataSink.tryEmitNext(datum); } // 发送完成后必须调用complete,让Flux知道没有更多数据 dataSink.tryEmitComplete(); return map; });
接下来,要实现非阻塞地等待Flux和Mono都完成,最终返回Mono的结果,可以利用flatMap操作符串联逻辑:
// 最终返回的非阻塞Mono Mono<HashSet<Long>> resultMono = mono.flatMap(map -> // 等待Flux消费完所有数据并完成,然后返回原Mono的结果map flux.then(Mono.just(map)) );
逻辑说明:
- 当订阅
resultMono时,首先执行mono的Supplier逻辑:发送所有数据到Sink,并完成Sink。 flux会在parallel调度器上消费所有数据,直到Sink发出完成信号。flux.then(Mono.just(map))会等待Flux完全结束后,直接返回之前mono生成的map,避免重复执行发送逻辑(因为mono是冷发布者,重复订阅会重新跑一遍发送流程)。
为什么不推荐你提到的其他方案:
- 阻塞的
block():会阻塞当前线程,违背响应式编程非阻塞、异步的核心原则,尤其在需要返回Mono给客户端的场景下,阻塞会破坏整个响应式链。 Flux.merge:它会把多个流的元素合并成一个序列,但无法直接提取原Mono的结果,还需要额外的逻辑去过滤和收集,复杂度更高,不如flatMap直接高效。
内容的提问来源于stack exchange,提问作者veenatic
相关产品推荐
相关产品推荐

