使用CompletableFuture.allOf仍响应缓慢,求异步调用优化方案
问题分析与优化方案
你的代码核心问题在于串行处理每个元素的异步任务,加上不合理的阻塞调用,导致元素数量增加时响应时间线性恶化。具体问题和优化方法如下:
一、当前实现的核心问题
- 元素级串行执行:外层for循环中,每个
SummaryList都要等待allOf().join()完成才会处理下一个,相当于5个元素就串行执行5轮异步任务,25个元素就是25轮,时间直接随元素数量翻倍增长。 - 异步线程内阻塞:在
runAsync的lambda里调用.get(),把异步调用硬生生变成同步阻塞,完全浪费了CompletableFuture的异步能力,还占用了线程池资源。 - 依赖默认线程池:
runAsync()默认使用ForkJoinPool.commonPool(),线程数有限(一般等于CPU核心数),大量任务提交后会排队等待,进一步拉长响应时间。 - 异常吞灭:catch块直接吞掉异常,排查问题时无法定位错误根源。
二、优化方案
1. 批量提交所有异步任务,统一等待完成
把所有元素的两个服务调用任务全部提交,最后用CompletableFuture.allOf()等待所有任务完成,实现全并行处理,避免元素间的串行等待。
2. 移除异步块内的阻塞调用,利用链式API
既然invokeFirstService和invokeSecondService本身返回CompletableFuture<Void>,直接用链式API处理结果,不需要额外包runAsync再调用.get()。如果服务可以修改,建议让它们返回带响应结果的CompletableFuture<FirstResponseBean>,逻辑会更简洁。
3. 使用自定义线程池控制并发
根据下游服务的承载能力,自定义线程池的核心线程数和最大线程数,避免并发过高压垮服务,同时避免默认线程池的资源竞争。
4. 合理处理异常
不要吞灭异常,至少记录日志,或者将异常传递出去,方便后续排查问题。
优化后的代码示例
首先定义全局复用的自定义线程池:
// 根据下游服务承载能力调整参数,比如核心线程数10,最大20 private static final ExecutorService CUSTOM_EXECUTOR = new ThreadPoolExecutor( 10, 20, 60L, TimeUnit.SECONDS, new LinkedBlockingQueue<>(100), new ThreadFactoryBuilder().setNameFormat("summary-service-pool-%d").build(), new ThreadPoolExecutor.CallerRunsPolicy() // 拒绝策略可根据业务调整 );
然后重构业务逻辑:
SummaryList s1=new SummaryList("101","account1","",""); SummaryList s2=new SummaryList("102","account2","",""); SummaryList s3=new SummaryList("103","account3","",""); SummaryList s4=new SummaryList("104","account4","",""); SummaryList s5=new SummaryList("105","account5","",""); List<SummaryList> list=Arrays.asList(s1,s2,s3,s4,s5); // 收集所有异步任务 List<CompletableFuture<Void>> allFutures = new ArrayList<>(); for (SummaryList summary : list) { FirstRequestBean firstRequest = new FirstRequestBean(); firstRequest.setAccount(summary.getAccount()); FirstResponseBean firstResponse = new FirstResponseBean(); // 处理第一个服务调用,链式处理结果 CompletableFuture<Void> firstFuture = invokeFirstService(firstRequest, firstResponse) .thenRun(() -> { if (!CollectionUtils.isEmpty(firstResponse.getResponse())) { summary.setName(firstResponse.getResponse().get(0).getName()); } }) .exceptionally(e -> { // 记录异常日志 log.error("调用invokeFirstService失败,account: {}", summary.getAccount(), e); return null; }) .asyncOn(CUSTOM_EXECUTOR); // 指定自定义线程池 SecondRequestBean secondRequest = new SecondRequestBean(); secondRequest.setAccount(summary.getAccount()); SecondResponseBean secondResponse = new SecondResponseBean(); // 处理第二个服务调用 CompletableFuture<Void> secondFuture = invokeSecondService(secondRequest, secondResponse) .thenRun(() -> { if (!CollectionUtils.isEmpty(secondResponse.getResponse())) { summary.setType(secondResponse.getResponse().get(0).getType()); } }) .exceptionally(e -> { log.error("调用invokeSecondService失败,account: {}", summary.getAccount(), e); return null; }) .asyncOn(CUSTOM_EXECUTOR); // 将当前元素的两个任务加入总列表 allFutures.add(firstFuture); allFutures.add(secondFuture); } // 等待所有任务完成 CompletableFuture.allOf(allFutures.toArray(new CompletableFuture[0])).join(); // 若线程池为单次使用则关闭,全局复用则无需此步骤 // CUSTOM_EXECUTOR.shutdown();
如果可以修改服务接口,让其返回带结果的Future,代码会更简洁:
// 假设服务修改为返回CompletableFuture<FirstResponseBean> CompletableFuture<Void> firstFuture = invokeFirstService(firstRequest) .thenAccept(response -> { if (!CollectionUtils.isEmpty(response.getResponse())) { summary.setName(response.getResponse().get(0).getName()); } }) .exceptionally(e -> { log.error("调用invokeFirstService失败,account: {}", summary.getAccount(), e); return null; }) .asyncOn(CUSTOM_EXECUTOR);
额外优化建议
- 请求合并:如果下游服务支持批量查询(比如传入多个account一次性返回结果),可以把所有account收集起来批量调用,减少IO次数,这是最有效的优化方式。
- 超时控制:给每个异步任务添加超时时间,避免单个慢请求拖垮整个流程,比如使用
orTimeout:firstFuture.orTimeout(2, TimeUnit.SECONDS); - 线程池复用:不要每次请求都创建新线程池,复用全局线程池,减少线程创建销毁的开销。
内容的提问来源于stack exchange,提问作者Tejesh Chegondi
相关产品推荐
相关产品推荐

