SpringBoot2.1.3 @Async配合CompletableFuture未在独立线程执行问题
问题结论
这个现象完全不符合并行异步调用的预期,你的代码存在两个核心错误,导致@Async异步逻辑完全失效,四个统计方法全部在Tomcat请求主线程里串行执行,根本没有并行能力。
错误原因
- 第一,Spring @Async注解基于AOP动态代理实现,只有通过Spring容器管理的代理对象调用带@Async的方法时,异步切面才会生效。你当前在
OrderStatsService类的compute()方法里直接调用同类下的四个getStatsX()方法,走的是this原生对象引用,完全绕过了代理对象,所以注解等于没加,方法直接同步执行。你日志里四个方法的线程名全是http-nio-8080-exec-1、getStats3比其他方法晚启动200多秒,就是串行同步执行的直接证据。 - 第二,你的线程池配置存在严重缺陷:初始化
ThreadPoolTaskExecutor后没有设置任何参数,默认核心线程数为1、队列无界、最大线程数为整数最大值,就算@Async生效,核心线程只有1个也没法同时跑4个任务,还可能因为无界队列引发OOM。
修复方案
你可以二选一实现,推荐第二种手动提交CompletableFuture的方案,完全绕开AOP代理的坑,逻辑更可控:
方案1:拆分异步方法到独立Bean,走AOP代理
把四个带@Async的统计方法抽到单独的Spring管理的Bean中,在OrderStatsService里注入这个Bean再调用,确保调用经过代理对象。
修正后的线程池配置
@EnableAsync @Configuration public class Config { @Bean(name = "threadPoolTaskExecutor") public Executor threadPoolTaskExecutor() { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); // 内存密集型任务核心线程数设为CPU核心数+1 int coreSize = Runtime.getRuntime().availableProcessors() + 1; executor.setCorePoolSize(coreSize); executor.setMaxPoolSize(coreSize * 2); executor.setQueueCapacity(100); executor.setThreadNamePrefix("stats-calc-"); executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy()); executor.initialize(); return executor; } }
独立的异步任务类
@Service public class StatsTaskService { @Async("threadPoolTaskExecutor") public CompletableFuture<Long> getStats1(List<Order> orders) { log.info("getStats1 started at " + System.currentTimeMillis() + ", Current Thread Name: " + Thread.currentThread().getName() + ", Current Thread ID: " + Thread.currentThread().getId()); // 统计1计算逻辑 return CompletableFuture.completedFuture(/*计算结果*/); } @Async("threadPoolTaskExecutor") public CompletableFuture<List<String>> getStats2(List<Order> orders) { log.info("getStats2 started at " + System.currentTimeMillis() + ", Current Thread Name: " + Thread.currentThread().getName() + ", Current Thread ID: " + Thread.currentThread().getId()); // 统计2计算逻辑 return CompletableFuture.completedFuture(/*计算结果*/); } @Async("threadPoolTaskExecutor") public CompletableFuture<Double> getStats3(List<Order> orders) { log.info("getStats3 started at " + System.currentTimeMillis() + ", Current Thread Name: " + Thread.currentThread().getName() + ", Current Thread ID: " + Thread.currentThread().getId()); // 统计3计算逻辑 return CompletableFuture.completedFuture(/*计算结果*/); } @Async("threadPoolTaskExecutor") public CompletableFuture<Map<String,String>> getStats4(List<Order> orders) { log.info("getStats4 started at " + System.currentTimeMillis() + ", Current Thread Name: " + Thread.currentThread().getName() + ", Current Thread ID: " + Thread.currentThread().getId()); // 统计4计算逻辑 return CompletableFuture.completedFuture(/*计算结果*/); } }
修正后的业务服务类
@Service class OrderStatsService{ @Autowired private StatsTaskService statsTaskService; public CumulativeStats compute() { log.info("CumulativeResult compute started at " + System.currentTimeMillis() + ", Current Thread Name: " + Thread.currentThread().getName() + ", Current Thread ID: " + Thread.currentThread().getId()); List<Order> orders = getOrders();// 拉取约10w条订单数据 CumulativeResult cumulativeResult = new CumulativeResult(); // 通过注入的代理对象调用异步方法 CompletableFuture<Long> stats1 = statsTaskService.getStats1(orders); CompletableFuture<List<String>> stats2 = statsTaskService.getStats2(orders); CompletableFuture<Double> stats3 = statsTaskService.getStats3(orders); CompletableFuture<Map<String,String>> stats4 = statsTaskService.getStats4(orders); // 等待所有任务执行完成后组装结果,可按需加超时、异常处理 CompletableFuture.allOf(stats1, stats2, stats3, stats4).join(); cumulativeResult.setStats1(stats1.join()); cumulativeResult.setStats2(stats2.join()); cumulativeResult.setStats3(stats3.join()); cumulativeResult.setStats4(stats4.join()); return cumulativeResult; } }
方案2:不用@Async,手动提交任务到线程池
直接通过CompletableFuture.supplyAsync()把任务提交到自定义线程池,完全不需要依赖AOP,不会出现内部调用失效的问题,代码逻辑更直观:
@Service class OrderStatsService{ @Autowired private Executor threadPoolTaskExecutor; public CumulativeStats compute() { log.info("CumulativeResult compute started at " + System.currentTimeMillis() + ", Current Thread Name: " + Thread.currentThread().getName() + ", Current Thread ID: " + Thread.currentThread().getId()); List<Order> orders = getOrders();// 拉取约10w条订单数据 CumulativeResult cumulativeResult = new CumulativeResult(); // 直接提交任务到线程池 CompletableFuture<Long> stats1 = CompletableFuture.supplyAsync(() -> { log.info("getStats1 started at " + System.currentTimeMillis() + ", Current Thread Name: " + Thread.currentThread().getName() + ", Current Thread ID: " + Thread.currentThread().getId()); // 统计1计算逻辑 return 0L; }, threadPoolTaskExecutor); CompletableFuture<List<String>> stats2 = CompletableFuture.supplyAsync(() -> { log.info("getStats2 started at " + System.currentTimeMillis() + ", Current Thread Name: " + Thread.currentThread().getName() + ", Current Thread ID: " + Thread.currentThread().getId()); // 统计2计算逻辑 return Collections.emptyList(); }, threadPoolTaskExecutor); CompletableFuture<Double> stats3 = CompletableFuture.supplyAsync(() -> { log.info("getStats3 started at " + System.currentTimeMillis() + ", Current Thread Name: " + Thread.currentThread().getName() + ", Current Thread ID: " + Thread.currentThread().getId()); // 统计3计算逻辑 return 0.0; }, threadPoolTaskExecutor); CompletableFuture<Map<String,String>> stats4 = CompletableFuture.supplyAsync(() -> { log.info("getStats4 started at " + System.currentTimeMillis() + ", Current Thread Name: " + Thread.currentThread().getName() + ", Current Thread ID: " + Thread.currentThread().getId()); // 统计4计算逻辑 return Collections.emptyMap(); }, threadPoolTaskExecutor); CompletableFuture.allOf(stats1, stats2, stats3, stats4).join(); cumulativeResult.setStats1(stats1.join()); cumulativeResult.setStats2(stats2.join()); cumulativeResult.setStats3(stats3.join()); cumulativeResult.setStats4(stats4.join()); return cumulativeResult; } }
验证标准
修复后查看日志,四个统计方法的线程名应为你配置的stats-calc-前缀的工作线程,线程ID和主线程不同,四个方法的启动时间差在几毫秒内,总执行时间约等于四个任务中最慢的那个的耗时,而非四个任务耗时总和。
内容的提问来源于stack exchange,提问作者OTUser
相关产品推荐
相关产品推荐

