如何防止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
相关产品推荐
相关产品推荐

