Java使用CompletableFuture.allOf启动多异步任务却顺序运行问题
问题原因
@Async注解未生效,Spring的@Async基于动态代理实现,以下两种情况会导致注解失效:- 启动类/配置类未添加
@EnableAsync注解,未开启Spring异步支持 fetchData方法和调用它的fetchDataForAllClients方法在同一个类中,同类调用不走代理对象,注解不生效
- 启动类/配置类未添加
- 你当前
fetchData方法的返回值是手动创建的CompletableFuture.completedFuture(counter),该返回值是已完成状态的Future,方法内的所有逻辑都会同步执行完毕才会返回,就算@Async生效也会出现串行执行的情况 - 额外问题:你代码里的
counter++是非线程安全的,多线程环境下会出现计数错误,需要替换为AtomicInteger类型
解决方法
推荐你直接用CompletableFuture.supplyAsync自定义线程池实现,比@Async更可控,还能方便的添加你需要的超时控制逻辑:
第一步:配置采集专用线程池
@Configuration public class FetchThreadPoolConfig { @Bean public Executor fetchDataThreadPool() { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); // 线程参数可根据实际业务调整 executor.setCorePoolSize(10); executor.setMaxPoolSize(20); executor.setQueueCapacity(100); executor.setThreadNamePrefix("fetch-data-"); executor.initialize(); return executor; } }
第二步:调整业务方法
将fetchData抽为独立的Bean方法,去掉@Async注解,修改为普通业务方法:
// 抽为独立的Service Bean @Service public class FetchDataService { // 将原来的counter替换为原子类 private final AtomicInteger counter = new AtomicInteger(0); @Resource private MyApiService myApiService; @Resource private FetchStatsService fetchStatsService; public Integer fetchData(final String date, final Integer clientId) { int currentCount = counter.incrementAndGet(); System.out.println(currentCount + ". FetchDataThread Started for "+ clientId + " at " + LocalDateTime.now()); boolean failed = false; String errorMsg = null; try { myApiService.fetchDataForClient(clientId, date, date); } catch (MyApiException exception) { failed = true; errorMsg = exception.getMessage(); } catch (Exception e) { failed = true; errorMsg = "执行异常:" + e.getMessage(); } // 正常执行/业务异常都直接落库 fetchStatsService.createFetchStats(clientId, date, failed, errorMsg); return currentCount; } }
第三步:调整调用逻辑,添加超时控制
@Service public class ClientFetchService { @Resource private Executor fetchDataThreadPool; @Resource private FetchDataService fetchDataService; @Resource private FetchStatsService fetchStatsService; private static final Logger LOGGER = LoggerFactory.getLogger(ClientFetchService.class); public void fetchDataForAllClients() { String previousDate = DateUtils.getPreviousDate(); List<Integer> clientIdList = PropertiesUtil.getClientIdList(); CompletableFuture.allOf(clientIdList.stream() .map(clientId -> CompletableFuture.supplyAsync(() -> fetchDataService.fetchData(previousDate, clientId), fetchDataThreadPool) // 添加超时控制,示例为30秒,可按需调整 .orTimeout(30, TimeUnit.SECONDS) .exceptionally(e -> { LOGGER.error("客户端{}采集数据失败", clientId, e); // 超时场景单独落库 if (e instanceof TimeoutException) { fetchStatsService.createFetchStats(clientId, previousDate, true, "采集超时"); } return null; }) .thenAcceptAsync(s -> System.out.println(s + ". FetchDataThread Finished for " + clientId + " at " + LocalDateTime.now()))) .toArray(CompletableFuture<?>[]::new)) .join(); } }
如果你要继续使用@Async方案,需满足以下要求:
- 启动类/配置类添加
@EnableAsync注解 fetchData方法必须和调用方法分属不同的Bean,避免同类调用- 方法直接返回业务结果,不要手动创建
CompletableFuture.completedFuture,Spring会自动帮你包装为异步CompletableFuture:
@Async public Integer fetchData(final String date, final Integer clientId) { // 业务逻辑和上面一致 }
内容的提问来源于stack exchange,提问作者AL̲̳I
相关产品推荐
相关产品推荐

