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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.26 08:57:00