使用CompletableFuture并行调用Rest接口抛出CompletionException问题咨询
问题根因
这个问题的核心诱因是默认RestTemplate的底层HTTP实现并发能力缺陷,具体分析如下:
- Spring RestTemplate默认使用JDK自带的
HttpURLConnection处理HTTP请求,该实现的连接复用逻辑在多线程并发调用场景下存在已知bug,会出现请求头、请求体写入错乱的情况,导致被调用方无法正常解析请求参数,串行调用时不存在资源竞争因此完全正常。 - 你看到的
CompletionException是异步任务执行异常的包装类,其内部嵌套的根异常就是HTTP请求发送失败对应的异常栈。 - 额外风险点:代码中直接使用
CompletableFuture.supplyAsync()默认的全局ForkJoin公共线程池,一旦公共线程池被其他业务逻辑占满,会导致请求长时间阻塞甚至超时,也可能衍生出各类并发问题。
修复方案
1. 替换RestTemplate的底层HTTP客户端
推荐替换为并发支持更成熟的OkHttp或者Apache HttpComponents,以OkHttp为例:
首先引入依赖(Maven示例):
<dependency> <groupId>com.squareup.okhttp3</groupId> <artifactId>okhttp</artifactId> <version>4.10.0</version> </dependency> <dependency> <groupId>org.springframework</groupId> <artifactId>spring-web</artifactId> <version>替换为你项目使用的Spring版本</version> </dependency>
然后配置RestTemplate使用OkHttp客户端:
@Bean public RestTemplate restTemplate() { OkHttp3ClientHttpRequestFactory factory = new OkHttp3ClientHttpRequestFactory(); // 按需配置超时时间,单位毫秒 factory.setConnectTimeout(3000); factory.setReadTimeout(5000); factory.setWriteTimeout(3000); return new RestTemplate(factory); }
2. 自定义业务线程池供CompletableFuture使用
避免使用全局公共线程池,配置独立的异步调用线程池:
@Bean public ExecutorService asyncHttpExecutor() { return new ThreadPoolExecutor( 10, // 核心线程数,可按业务量级调整 50, // 最大线程数,可按业务量级调整 60L, TimeUnit.SECONDS, new LinkedBlockingQueue<>(100), // 队列容量,可按业务量级调整 new ThreadFactoryBuilder().setNameFormat("async-http-%d").build() ); }
3. 优化业务代码逻辑
增加异常处理、超时控制,避免无限等待:
public HashMap<String,Double> getResp(String requestJson, double defaultScore, String anomalySelfInclUrl,String anomalySelfExclUrl){ CompletableFuture<Double> f1 = getHelper(requestJson,anomalySelfInclUrl, defaultScore); CompletableFuture<Double> f2 = getHelper(requestJson,anomalySelfExclUrl, defaultScore); HashMap<String,Double> respMap= new HashMap<>(); try { // 增加超时控制,避免无限阻塞,超时时间可按需调整 CompletableFuture.allOf(f1,f2).get(6, TimeUnit.SECONDS); respMap.put(Constants.selfInc, f1.get()); respMap.put(Constants.selfExcl, f2.get()); } catch (Exception e) { // 替换为你项目的异常处理逻辑 e.printStackTrace(); } return respMap; } public CompletableFuture<Double> getHelper(String requestJson, String url, double defaultScore) { // 传入自定义的异步HTTP线程池 return CompletableFuture.supplyAsync(() -> { HttpHeaders headers = new HttpHeaders(); headers.setContentType(MediaType.APPLICATION_JSON); HttpEntity<String> entity = new HttpEntity<>(requestJson, headers); String resp = restTemplate.postForObject(url, entity, String.class); JsonObject json = new JsonParser().parse(resp).getAsJsonObject(); return json.get("score").getAsDouble(); }, asyncHttpExecutor) // 单个请求失败返回默认值,不会阻断整个流程 .exceptionally(e -> { e.printStackTrace(); return defaultScore; }); }
内容的提问来源于stack exchange,提问作者Ayan Biswas
相关产品推荐
相关产品推荐

