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

Project Reactor实践:2个并行Flux REST请求成功后取消剩余请求

基于Project Reactor(Spring WebFlux)的实现方案

核心思路

利用Reactor的流式取消机制,并行发起所有请求,统计成功返回的结果数,拿到2个成功结果后主动取消上游所有未完成的请求,自动终止剩余API调用。


实现步骤

1. 定义单个API请求方法

首先基于WebClient封装单个REST请求,添加错误处理,避免单个请求失败打断整个主流程:

import org.springframework.web.reactive.function.client.WebClient;
import reactor.core.publisher.Mono;
import java.util.Objects;

// 初始化全局WebClient
WebClient webClient = WebClient.create("https://your-api-base-url.com");

// 单个API请求封装:请求成功返回结果,失败返回空Mono,不计入成功计数
private Mono<ApiResponse> callSingleApi(String apiPath) {
    return webClient.get()
            .uri(apiPath)
            .retrieve()
            .bodyToMono(ApiResponse.class)
            // 捕获所有请求错误,静默转为空,不中断主流程
            .onErrorResume(ex -> Mono.empty());
}

注:ApiResponse为自定义的接口返回值实体类,可根据业务替换为String或其他类型。

2. 并行请求+按条件终止逻辑

用Flux.merge并行订阅所有请求,结合take操作符实现目标逻辑:

import reactor.core.publisher.Flux;
import java.util.List;
import java.util.stream.Collectors;

public Flux<ApiResponse> runApisWithTwoSuccessThreshold(List<String> apiPathList) {
    // 把所有接口路径转为请求Mono
    List<Mono<ApiResponse>> allRequests = apiPathList.stream()
            .map(this::callSingleApi)
            .collect(Collectors.toList());

    return Flux.merge(allRequests) // 同时订阅所有请求,并行执行
            .filter(Objects::nonNull) // 过滤掉请求失败返回的空值,仅保留成功结果
            .take(2); // 收到2个成功结果后,立即向上游发送取消信号,终止所有未完成请求
}

如果需要把2个结果合并为集合返回,可在最后添加.collectList(),将Flux<ApiResponse>转为Mono<List<ApiResponse>>。


关键逻辑说明

  • 并行执行的实现:Flux.merge会同时订阅所有传入的Mono,和顺序执行的concat、按元素并行处理的parallel操作符不同,这个方案天然适合多个独立HTTP请求的并行发起场景,不需要额外配置线程池。
  • 自动终止的原理:take(2)操作符在收到指定数量的元素后,会主动向上游传播取消信号,Spring WebClient的底层已经适配了Reactor的取消机制,会自动终止对应未完成的HTTP请求、释放连接资源,不需要手动处理。
  • 错误兼容:单个请求出错会转为空Mono,不会触发整个流的异常,也不会占用成功计数,只有真正响应成功的请求才会被统计。

内容的提问来源于stack exchange,提问作者Bit Wise

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.29 12:54:04