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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.26 17:01:12