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
相关产品推荐
相关产品推荐

