Spring Reactor如何获取首个信号且不取消剩余Mono/Flux执行?
解决方案
核心思路是先把两个Mono转为不受下游取消影响的热发布序列,再取首个返回结果做优先处理,这样即使下游拿到首个结果后不再订阅,两个原始任务都会完整执行。
实现逻辑说明
- 给每个原始Mono添加
.cache()操作符,把冷序列转为热序列,缓存执行结果同时不会因为下游取消中断执行 - 用
Mono.firstWithSignal包装两个处理后的热Mono,拿到率先返回的结果执行优先处理逻辑 - 如果需要等待两个任务都全部执行完成再结束整个流程,可以再加一个
Mono.when来监听两个任务的完成信号
代码示例
// 两个原始异步任务示例 Mono<String> task1 = Mono.delay(Duration.ofSeconds(1)) .map(l -> "task1返回结果") .doOnTerminate(() -> System.out.println("task1执行完成")); Mono<String> task2 = Mono.delay(Duration.ofSeconds(3)) .map(l -> "task2返回结果") .doOnTerminate(() -> System.out.println("task2执行完成")); // 转为热序列,避免被下游取消中断执行 Mono<String> hotTask1 = task1.cache(); Mono<String> hotTask2 = task2.cache(); // 处理率先返回的结果 Mono<String> firstResultProcessor = Mono.firstWithSignal(hotTask1, hotTask2) .doOnNext(result -> System.out.println("拿到率先返回的结果:" + result + ",启动优先处理流程")); // 如果需要等两个任务都执行完再结束主流程,合并完成信号即可 Mono<Void> fullProcess = firstResultProcessor.then(Mono.when(hotTask1, hotTask2)); // 触发执行 fullProcess.block();
注意事项
- 如果不需要缓存任务执行结果,仅需保证任务不受取消影响跑完,可以把
cache()替换为publish().autoConnect(0),效果一致 - 如果不需要等待两个任务都执行完成再结束主流程,单独订阅
firstResultProcessor即可,两个后台任务仍会自行执行完毕
内容的提问来源于stack exchange,提问作者Frank Zhong
相关产品推荐
相关产品推荐

