如何在不修改CompletableFuture.supplyAsync的前提下拆分大请求批量调用REST API
如何在不修改CompletableFuture.supplyAsync逻辑的前提下拆分大请求批量调用REST API?
问题描述
我原本有一段代码用于发送包含数据列表的requestObj到REST API,但当requestObj过大时需要拆分成多个小列表,多次调用该API,但不想修改CompletableFuture.supplyAsync内部的逻辑:
ExecutorService executorService = Executors.newCachedThreadPool(); CompletableFuture.supplyAsync(() -> this.restTemplate.exchange(join(restAPI URL), POST, new HttpEntity<>(requestObj.toString(), httpHeaders), String.class), executorService) .thenAcceptAsync(response -> deleteRow(response, deleteCount), executorService) .exceptionally(error -> { logger.error("Error :: " + error); return null; });
解决方案
这问题我之前也碰到过,核心就是把拆分后的每个小请求复用你原有的异步调用链,完全不需要改动supplyAsync里的REST调用逻辑。具体步骤如下:
1. 拆分大请求为多个小请求
首先你需要把原来的大requestObj拆分成多个包含子列表的小requestObj。这里假设你的requestObj有获取数据列表的方法,以及通过子列表构造新实例的方式:
// 自定义拆分方法,根据API的请求大小限制调整子列表阈值 private List<RequestObj> splitLargeRequest(RequestObj originalRequest) { List<DataItem> originalDataList = originalRequest.getDataList(); // 这里用Guava的Lists.partition做拆分,也可以自己实现逻辑 List<List<DataItem>> splitDataLists = Lists.partition(originalDataList, 100); // 将每个子列表转换成新的RequestObj实例 return splitDataLists.stream() .map(subList -> new RequestObj(subList)) // 替换成你实际的构造方式 .collect(Collectors.toList()); }
2. 封装原有的异步调用逻辑
把你原来的CompletableFuture调用链封装成一个可复用的方法,这样每个小请求都能直接复用这段逻辑,完全保留原有的supplyAsync、thenAcceptAsync和异常处理:
// 封装原异步逻辑,参数按需调整 private CompletableFuture<Void> processSingleRequest( RequestObj requestObj, ExecutorService executorService, HttpHeaders httpHeaders, String restApiUrl, int deleteCount) { // 这里完全照搬你原来的代码,没有任何修改 return CompletableFuture.supplyAsync(() -> this.restTemplate.exchange(restApiUrl, HttpMethod.POST, new HttpEntity<>(requestObj.toString(), httpHeaders), String.class), executorService) .thenAcceptAsync(response -> deleteRow(response, deleteCount), executorService) .exceptionally(error -> { logger.error("Error processing request batch :: " + error.getMessage(), error); return null; }); }
3. 批量执行所有异步任务
最后,拆分请求后批量生成异步任务,并用CompletableFuture.allOf()统一管理所有任务的完成状态:
// 你的原始大请求对象 RequestObj bigRequestObj = ...; // REST API地址、HttpHeaders等参数 String restApiUrl = join(restAPI URL); HttpHeaders httpHeaders = ...; int deleteCount = ...; // 拆分大请求 List<RequestObj> splitRequests = splitLargeRequest(bigRequestObj); // 生成所有异步任务 List<CompletableFuture<Void>> taskFutures = splitRequests.stream() .map(request -> processSingleRequest(request, executorService, httpHeaders, restApiUrl, deleteCount)) .collect(Collectors.toList()); // 等待所有批量任务完成(如果不需要阻塞主线程,可以用thenRun做后续处理) CompletableFuture.allOf(taskFutures.toArray(new CompletableFuture[0])).join(); // 非阻塞写法示例: // CompletableFuture.allOf(taskFutures.toArray(new CompletableFuture[0])) // .thenRun(() -> logger.info("All batch requests have finished processing"));
注意事项
- 子列表大小:根据你的REST API允许的最大请求体大小来调整拆分的阈值,避免小请求仍然超出限制。
- 线程池配置:
newCachedThreadPool会动态创建线程,批量调用时可能创建过多线程导致资源耗尽,建议改用固定大小的线程池(比如Executors.newFixedThreadPool(10))更稳妥。 - 顺序需求:如果API调用需要保证顺序,可以用
thenCompose链式调用,但并行处理通常效率更高,适合无顺序要求的场景。
内容的提问来源于stack exchange,提问作者David
相关产品推荐
相关产品推荐

