Spring WebFlux如何实现请求合并批量调用慢API并拆分响应
Flux实现请求批量聚合+结果路由方案
核心思路
不要直接向缓冲队列投递裸查询参数,而是将每个入站请求的查询参数、对应的异步响应发射器绑定为上下文对象再入队;批量请求拿到全量结果后,通过绑定的发射器将对应子集结果发回原始请求,全程不会丢失请求上下文。bufferTimeout操作符完全匹配「队列容量达阈值/等待超时触发批量请求」的需求,不需要用复杂度更高的Window操作符。
实现步骤
- 定义请求上下文包装类,包含两个核心字段:当前请求携带的查询参数集合、当前请求专属的响应结果发射器。
- 构造线程安全的全局响应式缓冲队列,支持多请求并发提交,配置合理的背压策略避免服务被突增流量打崩。
- 对队列的Flux视图挂载
bufferTimeout操作符,直接配置批量大小阈值、超时时间两个核心参数。 - 批次触发时,先合并批次内所有查询参数并去重,再通过WebClient发起一次第三方批量请求拿到全量结果。
- 遍历批次内的所有请求上下文,从全量结果中筛出当前请求对应的参数子集结果,通过绑定的发射器将结果发回给原始请求。
- 补充异常处理逻辑:第三方API调用失败时,必须给批次内所有请求发送错误信号,避免原始请求无限等待挂死。
核心代码示例
import org.springframework.web.bind.annotation.GetMapping; import org.springframework.web.bind.annotation.RequestParam; import org.springframework.web.bind.annotation.RestController; import org.springframework.web.reactive.function.client.WebClient; import reactor.core.publisher.Mono; import reactor.core.publisher.Sinks; import jakarta.annotation.PostConstruct; import java.time.Duration; import java.util.Map; import java.util.Set; import java.util.stream.Collectors; import org.springframework.core.ParameterizedTypeReference; @RestController public class ShippingBatchController { private final WebClient webClient = WebClient.create(); // 全局线程安全缓冲队列 private final Sinks.Many<ShippingRequestContext> requestBuffer = Sinks.many() .multicast() .onBackpressureBuffer(1000); // 队列最大长度按需调整 // 单请求上下文:绑定查询参数和响应发射器 record ShippingRequestContext( Set<String> queryCodes, Sinks.One<Map<String, Integer>> responseSink ) {} @PostConstruct public void initBatchProcessor() { requestBuffer.asFlux() // 核心配置:凑满5个参数 或 等待5秒 即触发批量请求 .bufferTimeout(5, Duration.ofSeconds(5)) .flatMap(batch -> { // 合并所有查询参数并去重,减少冗余参数 Set<String> allQueryCodes = batch.stream() .flatMap(ctx -> ctx.queryCodes().stream()) .collect(Collectors.toSet()); // 调用第三方慢API return webClient.get() .uri("https://slowapi.com/shipping?q=" + String.join(",", allQueryCodes)) .retrieve() .bodyToMono(new ParameterizedTypeReference<Map<String, Integer>>() {}) // 异常时通知批次内所有请求 .onErrorResume(err -> { batch.forEach(ctx -> ctx.responseSink().tryEmitError(err)); return Mono.empty(); }) // 拆分全量结果,路由到对应请求 .doOnNext(fullResult -> { batch.forEach(ctx -> { Map<String, Integer> subResult = fullResult.entrySet().stream() .filter(entry -> ctx.queryCodes().contains(entry.getKey())) .collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue)); ctx.responseSink().tryEmitValue(subResult); }); }); }) .subscribe(); } @GetMapping("/shipping") public Mono<Map<String, Integer>> getShippingFee(@RequestParam("q") Set<String> queryCodes) { // 为当前请求创建独立的响应发射器 Sinks.One<Map<String, Integer>> responseSink = Sinks.one(); // 上下文入队 requestBuffer.tryEmitNext(new ShippingRequestContext(queryCodes, responseSink)); // 返回异步结果,等批量处理完成后自动返回对应数据 return responseSink.asMono(); } }
生产环境注意事项
- 不要直接缓冲裸查询参数:如果队列中只存储参数不绑定对应响应发射器,异步场景下完全无法将结果映射回原始请求,是该场景最常见的实现错误。
- 参数合并前必须去重:多个请求携带重复查询参数时,去重可以有效缩短第三方API的请求参数长度,进一步降低调用开销。
- 可以根据第三方API的性能上限调整批量大小、超时时间两个参数,在调用次数和请求延迟之间找平衡点。
- 缓冲队列需要配置最大长度和背压拒绝策略,避免流量突增时队列积压过多导致内存溢出。
内容的提问来源于stack exchange,提问作者de.la.ru
相关产品推荐
相关产品推荐

