如何提升CompletableFuture Join性能?Java批量REST调用优化求助
优化CompletableFuture并行调用REST API的性能方案
针对4000个用户创建请求耗时剧增的问题,核心原因是无限制并发引发的线程资源耗尽、连接开销过大、目标服务限流等问题,以下是具体优化手段:
1. 自定义线程池控制并发度
默认CompletableFuture使用的ForkJoinPool.commonPool()线程数仅为CPU核心数左右,对于IO密集型的HTTP调用,线程数不足会导致任务排队;但无限制开线程又会引发上下文切换开销。建议针对IO场景配置线程池:
// 自定义IO密集型线程池 ThreadPoolExecutor executor = new ThreadPoolExecutor( Runtime.getRuntime().availableProcessors() * 2, // 核心线程数,CPU核心*2 Runtime.getRuntime().availableProcessors() * 4, // 最大线程数 60L, TimeUnit.SECONDS, new LinkedBlockingQueue<>(100), // 任务缓冲队列 Executors.defaultThreadFactory(), new ThreadPoolExecutor.CallerRunsPolicy() // 队列满时让主线程兜底处理,避免丢任务 ); // 提交任务时指定自定义线程池 List<CompletableFuture<Object>> allUsersFuturesCalls = new ArrayList<>(); for (User user : listOfUsers) { allUsersFuturesCalls.add(CompletableFuture.supplyAsync(() -> { User userToCreate = userService.MapUser(user); return userService.createUser(userToCreate); }, executor)); } CompletableFuture.allOf(allUsersFuturesCalls.toArray(new CompletableFuture[0])).join(); executor.shutdown(); // 任务完成后关闭线程池
2. 给RestTemplate配置连接池
默认RestTemplate每次HTTP请求都会新建连接,4000次请求会产生大量TCP连接建立/销毁开销。改用HttpClient连接池复用连接:
// 配置HttpClient连接池 PoolingHttpClientConnectionManager connectionManager = new PoolingHttpClientConnectionManager(); connectionManager.setMaxTotal(200); // 最大总连接数 connectionManager.setDefaultMaxPerRoute(100); // 单个目标地址的最大连接数 CloseableHttpClient httpClient = HttpClientBuilder.create() .setConnectionManager(connectionManager) .setDefaultRequestConfig(RequestConfig.custom() .setConnectTimeout(5000) // 连接超时5秒 .setSocketTimeout(10000) // 读取超时10秒 .build()) .build(); // 替换RestTemplate的请求工厂 HttpComponentsClientHttpRequestFactory requestFactory = new HttpComponentsClientHttpRequestFactory(httpClient); RestTemplate restTemplate = new RestTemplate(requestFactory);
3. 分批处理+批量请求(API支持的话)
如果目标REST API提供批量创建接口,将4000个用户分成若干批次(比如每50个一批),大幅减少HTTP请求次数:
// 按50个一批拆分用户列表(可自己实现拆分逻辑,或用Guava的Lists.partition) List<List<User>> batches = new ArrayList<>(); for (int i = 0; i < listOfUsers.size(); i += 50) { batches.add(listOfUsers.subList(i, Math.min(i + 50, listOfUsers.size()))); } List<CompletableFuture<Object>> batchFutures = new ArrayList<>(); for (List<User> batch : batches) { batchFutures.add(CompletableFuture.supplyAsync(() -> { List<User> usersToCreate = batch.stream() .map(userService::MapUser) .collect(Collectors.toList()); return userService.batchCreateUsers(usersToCreate); // 调用批量创建接口 }, executor)); } CompletableFuture.allOf(batchFutures.toArray(new CompletableFuture[0])).join();
4. 增加超时控制,避免任务阻塞
给每个异步任务设置超时,防止个别慢请求拖慢整体进度:
allUsersFuturesCalls.add(CompletableFuture.supplyAsync(() -> { User userToCreate = userService.MapUser(user); return userService.createUser(userToCreate); }, executor).orTimeout(15, TimeUnit.SECONDS)); // 15秒超时,超时抛出TimeoutException
5. 分阶段处理,降低内存压力
一次性提交4000个CompletableFuture会占用大量内存,可分批次提交并等待完成后再处理下一批:
int batchSize = 500; for (int i = 0; i < listOfUsers.size(); i += batchSize) { List<User> subList = listOfUsers.subList(i, Math.min(i + batchSize, listOfUsers.size())); List<CompletableFuture<Object>> batchFutures = new ArrayList<>(); for (User user : subList) { batchFutures.add(CompletableFuture.supplyAsync(() -> { User userToCreate = userService.MapUser(user); return userService.createUser(userToCreate); }, executor)); } CompletableFuture.allOf(batchFutures.toArray(new CompletableFuture[0])).join(); }
6. 并行处理映射与API调用
如果MapUser是CPU密集型操作,可将其与HTTP调用并行处理,避免主线程阻塞:
// 先异步完成对象映射,再异步调用API allUsersFuturesCalls.add(CompletableFuture.supplyAsync(() -> userService.MapUser(user), executor) .thenCompose(userToCreate -> CompletableFuture.supplyAsync(() -> userService.createUser(userToCreate), executor)));
内容的提问来源于stack exchange,提问作者James
相关产品推荐
相关产品推荐

