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

如何防止Mono被取消?Reactor中返回快请求同时处理慢请求结果的方案

竞态请求不丢弃慢请求处理方案

问题本质

Flux.first() 操作符在获取到第一个完成的信号后,会主动取消对其余未完成Publisher的订阅。默认的冷Mono收到取消信号后会直接终止执行,导致慢请求的处理逻辑被丢弃。

Reactor 原生解决方案

核心思路是将两个请求转为热Publisher,只要触发订阅就会全程执行,不受下游取消操作影响,同时提前绑定每个请求的成功/失败处理逻辑:

import java.time.Duration;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;

// 改造后的核心代码
Mono<StatusMock> monoA = webClient.get()
  .uri("https://some.url.a")
  .retrieve()
  .bodyToMono(StatusMock.class)
  // 绑定请求A的独立处理逻辑,无论快慢都会执行
  .doOnSuccess(res -> {
    // 此处编写请求A成功后的处理逻辑,比如日志、落库等
    log.info("请求A执行完成,结果:{}", res.getStatus());
  })
  .doOnError(err -> {
    // 此处编写请求A失败后的处理逻辑
    log.error("请求A执行失败", err);
  })
  .subscribeOn(Schedulers.boundedElastic())
  // 核心:转为热Mono,只要被订阅一次就会执行到底,不受下游取消影响
  .cache();

Mono<StatusMock> monoB = webClient.get()
  .uri("https://some.url.b")
  .retrieve()
  .bodyToMono(StatusMock.class)
  .doOnSuccess(this::verifyBody)
  // 绑定请求B的独立处理逻辑
  .doOnSuccess(res -> {
    log.info("请求B执行完成,结果:{}", res.getStatus());
  })
  .doOnError(err -> {
    log.error("请求B执行失败", err);
  })
  .onErrorStop()
  .subscribeOn(Schedulers.boundedElastic())
  .cache();

// 合并两个请求,取第一个成功返回的非异常结果
StatusMock statusMock = Flux.merge(
    // 异常结果转为空,避免先报错的请求中断流程
    monoA.onErrorResume(e -> Mono.empty()),
    monoB.onErrorResume(e -> Mono.empty())
  )
  .next()
  // 增加超时避免两个请求都失败时永久阻塞
  .block(Duration.ofSeconds(10));

if (statusMock != null) {
  return statusMock.getStatus();
}
return "empty";

关键逻辑说明

  • cache() 操作符将冷Mono转为带缓存的热Mono,仅需一次订阅就会完整执行整个流程,下游的取消操作不会中断请求执行
  • 每个请求的处理逻辑提前绑定到对应Mono的生命周期中,无论请求是先完成还是后完成,都会在结束时触发对应处理逻辑
  • 合并请求时通过onErrorResume过滤异常结果,保证拿到的是最先成功的正常响应,而非最先完成的异常响应

备选方案:Java 原生 CompletableFuture 实现

如果不想依赖Reactor特性,可直接用JDK自带的CompletableFuture实现:

import java.util.concurrent.CompletableFuture;
import java.util.concurrent.Executor;

// 自行初始化线程池,也可复用Reactor的boundedElastic线程池
Executor executor = Schedulers.boundedElastic().asExecutor();

CompletableFuture<StatusMock> futureA = CompletableFuture.supplyAsync(() -> 
    webClient.get().uri("https://some.url.a")
      .retrieve()
      .bodyToMono(StatusMock.class)
      .block(), executor)
  // 绑定请求A的处理逻辑
  .whenComplete((res, err) -> {
    if (err != null) {
      log.error("请求A执行失败", err);
      return;
    }
    log.info("请求A执行完成,结果:{}", res.getStatus());
  });

CompletableFuture<StatusMock> futureB = CompletableFuture.supplyAsync(() -> {
    StatusMock res = webClient.get().uri("https://some.url.b")
      .retrieve()
      .bodyToMono(StatusMock.class)
      .block();
    verifyBody(res);
    return res;
  }, executor)
  .whenComplete((res, err) -> {
    if (err != null) {
      log.error("请求B执行失败", err);
      return;
    }
    log.info("请求B执行完成,结果:{}", res.getStatus());
  });

// 取第一个完成的正常结果
StatusMock statusMock = (StatusMock) CompletableFuture.anyOf(futureA, futureB)
  .exceptionally(e -> null)
  .join();

if (statusMock != null) {
  return statusMock.getStatus();
}
return "empty";

内容的提问来源于stack exchange,提问作者Marco Queiroz

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.05 08:39:01