Spring WebFlux中实现请求排队批量调用外部API方案咨询
批量请求外部API的订单服务实现方案
需求背景
现有OrderService直接调用外部API获取订单,外部API按请求计费且单次最多支持10个订单号。需要实现订单号排队逻辑:凑够10个订单号再发起批量请求,将结果准确返回给对应的原始请求;若队列中订单号超过10,取前10个处理,剩余留待下一批。
实现方案
基于Reactor的Sinks实现请求收集与批量分发,通过封装请求上下文跟踪每个原始请求对应的订单号,确保结果准确返回。
完整代码实现
@Component public class OrderService { private final WebClient webClient; private final Sinks.Many<OrderRequest> requestSink; // 封装单个请求的订单号与结果接收器 private static class OrderRequest { private final List<String> orderNumbers; private final MonoSink<Map<String, Order>> sink; OrderRequest(List<String> orderNumbers, MonoSink<Map<String, Order>> sink) { this.orderNumbers = orderNumbers; this.sink = sink; } } public OrderService(WebClient webClient) { this.webClient = webClient; // 创建带背压缓冲的多播Sink,用于收集所有订单请求 this.requestSink = Sinks.many().multicast().onBackpressureBuffer(); // 启动批量处理流程 startBatchProcessing(); } public Mono<List<Order>> getOrders(List<String> orderNumbers) { return Mono.create(sink -> { // 将当前请求封装后发送至Sink requestSink.emitNext(new OrderRequest(orderNumbers, sink), Sinks.EmitFailureHandler.FAIL_FAST); }).map(resultMap -> orderNumbers.stream() .map(resultMap::get) .collect(Collectors.toList())); } private void startBatchProcessing() { requestSink.asFlux() // 将每个请求的订单号拆分为单个条目,关联对应的请求接收器 .flatMapIterable(request -> request.orderNumbers.stream() .map(orderNum -> new AbstractMap.SimpleEntry<>(orderNum, request.sink))) // 凑够10个订单号组成一批 .buffer(10) // 串行处理每一批(避免并发调用外部API过多) .concatMap(this::processBatch) .subscribe(); } private Mono<Void> processBatch(List<Map.Entry<String, MonoSink<Map<String, Order>>>> batch) { // 提取当前批次的所有订单号 List<String> batchOrderNumbers = batch.stream() .map(Map.Entry::getKey) .collect(Collectors.toList()); // 调用外部API获取批量订单数据 return webClient.get() .uri(builder -> builder.queryParam("orderNumbers", String.join(",", batchOrderNumbers)).build()) .retrieve() .bodyToFlux(Order.class) .collectMap(Order::getOrderNumber) // 转换为订单号->订单的映射 .doOnSuccess(orderMap -> { // 按请求接收器分组,整理每个请求对应的结果 Map<MonoSink<Map<String, Order>>, Map<String, Order>> sinkResults = new HashMap<>(); batch.forEach(entry -> { String orderNum = entry.getKey(); MonoSink<Map<String, Order>> sink = entry.getValue(); sinkResults.computeIfAbsent(sink, k -> new HashMap<>()) .put(orderNum, orderMap.get(orderNum)); }); // 完成每个请求的结果返回 sinkResults.forEach((sink, result) -> sink.success(result)); }) .doOnError(error -> { // 出错时通知当前批次所有关联的请求 Set<MonoSink<Map<String, Order>>> affectedSinks = batch.stream() .map(Map.Entry::getValue) .collect(Collectors.toSet()); affectedSinks.forEach(sink -> sink.error(error)); }) .then(); // 表示批量处理完成 } }
核心逻辑说明
- 请求收集:每个
getOrders请求通过Mono.create生成响应流,同时将订单号和对应的MonoSink(结果控制器)封装后发送到Sink,实现请求的异步收集。 - 批量分组:将Sink转为Flux后,拆分每个请求的订单号为单个条目,通过
buffer(10)自动凑够10个订单号组成一批,不足10个时等待后续请求补充。 - 结果分发:调用外部API后,将响应转为订单号映射,按原始请求的
MonoSink分组整理结果,调用sink.success()将对应结果返回给发起请求的调用方;若API调用失败,所有关联请求都会收到错误信号。 - 兼容性:保持原方法返回类型
Mono<List<Order>>不变,无需修改上层调用逻辑。
内容的提问来源于stack exchange,提问作者orderlyfashion
相关产品推荐
相关产品推荐

