TaskExecutor 与异步 - @EnableAsync、ThreadPoolTaskExecutor 与 AsyncExecutionAspectSupport 源码
概述
Spring Framework 从 2.0 版本开始就提供了 TaskExecutor 抽象层,它将 Java 原生的线程池(java.util.concurrent.Executor)封装为 Spring 风格的接口,并结合 Spring 的声明式编程模型,通过 @Async 注解和 @EnableAsync 让开发者得以用极少的代码实现方法的异步执行。
本文从接口体系出发,逐步深入到注解驱动、异步配置、线程池参数调优、异常处理、上下文传递等核心主题,最后通过一个完整的异步导出 Excel 实战案例收尾。源码分析基于 Spring Framework 5.3.x。
1. TaskExecutor 接口体系
Spring 的异步执行抽象建立在 org.springframework.core.task 包下的一组接口之上,它们构成了从简单到复杂、从同步到异步的完整层级。
1.1 TaskExecutor —— 最顶层的抽象
TaskExecutor 是 Spring 中所有执行器的根接口,它定义了一个单一方法 execute(Runnable),语义等同于 java.util.concurrent.Executor.execute(Runnable)。
// org.springframework.core.task.TaskExecutor
@FunctionalInterface
public interface TaskExecutor extends Executor {
/**
* 执行给定的任务 {@code task}。
* 具体执行方式(同步/异步/单线程/线程池)由实现类决定。
*/
@Override
void execute(Runnable task);
}这个接口极度精简,但它统一了 Spring 容器内所有"执行任务"的行为——无论底层是简单的线程创建、线程池复用,还是消息队列委托,调用方看到的都是同一个 execute(Runnable) 契约。
1.2 AsyncListenableTaskExecutor —— 带回调的执行器
// org.springframework.core.task.AsyncListenableTaskExecutor
public interface AsyncListenableTaskExecutor extends AsyncTaskExecutor {
/**
* 提交任务并返回一个 ListenableFuture,调用方可以注册回调
* 以在任务完成或失败时获得通知。
*/
ListenableFuture<?> submitListenable(Runnable task);
}ListenableFuture 是 Spring 对 Future 的增强,它允许通过回调方式处理异步结果,无需阻塞等待。Spring 4.0 之后虽然有了 CompletableFuture,但 ListenableFuture 在部分遗留场景中仍有使用。
1.3 AsyncTaskExecutor —— 支持 Future 返回的执行器
// org.springframework.core.task.AsyncTaskExecutor
public interface AsyncTaskExecutor extends TaskExecutor {
/** 表示立即执行的常量 */
long TIMEOUT_IMMEDIATE = 0;
/** 表示无超时的常量 */
long TIMEOUT_INDEFINITE = Long.MAX_VALUE;
/**
* 提交一个 Callable 任务并返回 Future。
*/
<T> Future<T> submit(Callable<T> task);
/**
* 提交一个 Runnable 任务并返回 Future。
*/
Future<?> submit(Runnable task);
}1.4 SimpleAsyncTaskExecutor —— 简单实现
SimpleAsyncTaskExecutor 是 TaskExecutor 的一个简单实现,它不为线程做复用——每次 execute() 调用都会 创建一个新线程。
// org.springframework.core.task.SimpleAsyncTaskExecutor
public class SimpleAsyncTaskExecutor extends CustomizableThreadCreator
implements AsyncListenableTaskExecutor, Serializable {
private final Object concurrencyThrottle = new SimpleConcurrencyThrottle();
@Override
public void execute(Runnable task, long startTimeout) {
// 并发节流检查
concurrencyThrottle.beforeAccess();
doExecute(new ConcurrencyThrottlingRunnable(task));
}
protected void doExecute(Runnable task) {
// 每次调用都创建新线程
new Thread(task).start();
}
}注意
SimpleAsyncTaskExecutor 在生产环境中 不推荐使用。它不重用线程,高并发下会创建大量线程导致 OOM。它仅适用于测试或极端轻量级的场景。
1.5 SyncTaskExecutor —— 同步执行
SyncTaskExecutor 不是一个异步执行器,它 在当前线程中同步执行 提交的任务,相当于直接调用 task.run()。
// org.springframework.core.task.SyncTaskExecutor
public class SyncTaskExecutor implements TaskExecutor, Serializable {
@Override
public void execute(Runnable task) {
// 直接在当前线程同步执行
task.run();
}
}它在哪些场景下有用?当你需要"关闭异步"进行测试或调试时,可以将配置中的执行器替换为 SyncTaskExecutor,从而在不改业务代码的前提下让异步变成同步,方便调试。
1.6 ThreadPoolTaskExecutor —— 核心实现
ThreadPoolTaskExecutor 是 Spring 中 最核心、最常用 的执行器实现。它内部委托给 java.util.concurrent.ThreadPoolExecutor,并通过 Spring 生命周期管理(InitializingBean/DisposableBean)实现自动初始化和销毁。
// org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor
public class ThreadPoolTaskExecutor extends ExecutorConfigurationSupport
implements AsyncListenableTaskExecutor, SchedulingTaskExecutor {
// 核心线程数
private int corePoolSize = 1;
// 最大线程数
private int maxPoolSize = Integer.MAX_VALUE;
// 队列容量
private int queueCapacity = Integer.MAX_VALUE;
// 线程存活时间(秒)
private int keepAliveSeconds = 60;
// 是否允许核心线程超时(默认 false)
private boolean allowCoreThreadTimeOut = false;
// 线程名前缀
private String threadNamePrefix;
// 底层的 JDK ThreadPoolExecutor
private ThreadPoolExecutor threadPoolExecutor;
@Override
public void afterPropertiesSet() {
// Spring 容器初始化时调用,创建底层的 ThreadPoolExecutor
initialize();
}
/**
* 根据参数创建 java.util.concurrent.ThreadPoolExecutor
*/
protected ExecutorService initializeExecutor(
ThreadFactory threadFactory, RejectedExecutionHandler rejectedExecutionHandler) {
BlockingQueue<Runnable> queue = createQueue(this.queueCapacity);
ThreadPoolExecutor executor;
if (this.taskDecorator != null) {
// 如果配置了 TaskDecorator,则包装线程工厂
}
executor = new ThreadPoolExecutor(
this.corePoolSize,
this.maxPoolSize,
this.keepAliveSeconds, TimeUnit.SECONDS,
queue,
threadFactory,
rejectedExecutionHandler);
if (this.allowCoreThreadTimeOut) {
executor.allowCoreThreadTimeOut(true);
}
this.threadPoolExecutor = executor;
return executor;
}
/**
* 根据 queueCapacity 创建队列:
* - queueCapacity <= 0 → SynchronousQueue(无缓冲)
* - queueCapacity > 0 → LinkedBlockingQueue(有界队列)
*/
protected BlockingQueue<Runnable> createQueue(int queueCapacity) {
if (queueCapacity > 0) {
return new LinkedBlockingQueue<>(queueCapacity);
} else {
return new SynchronousQueue<>();
}
}
@Override
public void execute(Runnable task) {
ThreadPoolExecutor executor = getThreadPoolExecutor();
// 如果配置了 TaskDecorator,在提交前包装任务
if (this.taskDecorator != null) {
task = this.taskDecorator.decorate(task);
}
executor.execute(task);
}
@Override
public <T> Future<T> submit(Callable<T> task) {
ThreadPoolExecutor executor = getThreadPoolExecutor();
return executor.submit(task);
}
@Override
public ListenableFuture<?> submitListenable(Runnable task) {
ListenableFutureTask<Object> future = new ListenableFutureTask<>(task, null);
execute(future);
return future;
}
}关于 ThreadPoolTaskExecutor 的参数和行为细节,将在第 4 节中详细展开。
2. @EnableAsync 注解源码分析
@EnableAsync 是 Spring 异步功能的 总开关,它被定义在 org.springframework.scheduling.annotation 包中。
2.1 注解定义
// org.springframework.scheduling.annotation.EnableAsync
@Target(ElementType.TYPE)
@Retention(RetentionPolicy.RUNTIME)
@Documented
@Import(AsyncConfigurationSelector.class) // 关键!通过 @Import 引入配置选择器
public @interface EnableAsync {
/**
* 默认情况下,Spring 会搜索 @Async 注解标注的方法。
* 如果设置此属性,则可以搜索自定义的注解类型。
*/
Class<? extends Annotation> annotation() default Annotation.class;
/**
* 切面模式:
* - ADVICE_MODE_PROXY(默认):基于 JDK 动态代理或 CGLIB 代理
* - ADVICE_MODE_ASPECTJ:基于 AspectJ 编译期织入
*/
AdviceMode mode() default AdviceMode.PROXY;
/**
* 排序顺序,当存在多个切面时控制执行顺序。
*/
int order() default Ordered.LOWEST_PRECEDENCE;
/**
* 是否允许在同一个类内部通过"自调用"的方式
* 调用 @Async 方法(仅代理模式下有效)。
* 默认 false——即自调用不会触发异步。
*/
boolean proxyTargetClass() default false;
}2.2 @Import(AsyncConfigurationSelector.class)
@EnableAsync 通过 @Import(AsyncConfigurationSelector.class) 导入配置,这是 Spring 注解驱动的核心机制。
// org.springframework.scheduling.annotation.AsyncConfigurationSelector
public class AsyncConfigurationSelector extends AdviceModeImportSelector<EnableAsync> {
private static final String ASYNC_EXECUTION_ASPECT_CONFIGURATION_CLASS_NAME =
"org.springframework.scheduling.aspectj.AspectJAsyncConfiguration";
@Override
@Nullable
public String[] selectImports(AdviceMode adviceMode) {
switch (adviceMode) {
case PROXY:
// 默认模式:返回 ProxyAsyncConfiguration
return new String[]{ProxyAsyncConfiguration.class.getName()};
case ASPECTJ:
// AspectJ 模式:返回 AspectJAsyncConfiguration
return new String[]{ASYNC_EXECUTION_ASPECT_CONFIGURATION_CLASS_NAME};
default:
return null;
}
}
}2.3 ProxyAsyncConfiguration —— 代理模式下的异步配置
// org.springframework.scheduling.annotation.ProxyAsyncConfiguration
@Configuration
@Role(BeanDefinition.ROLE_INFRASTRUCTURE)
public class ProxyAsyncConfiguration extends AbstractAsyncConfiguration {
@Bean(name = TaskManagementConfigUtils.ASYNC_ANNOTATION_PROCESSOR_BEAN_NAME)
@Role(BeanDefinition.ROLE_INFRASTRUCTURE)
public AsyncAnnotationBeanPostProcessor asyncAdvisor() {
AsyncAnnotationBeanPostProcessor bpp = new AsyncAnnotationBeanPostProcessor();
// 设置自定义的异步执行器(如果有的话)
if (this.executor != null) {
bpp.configure(this.executor, this.exceptionHandler);
}
// 设置 @Async 注解类型
Class<? extends Annotation> customAsyncAnnotation = this.enableAsync.get("annotation");
if (customAsyncAnnotation != Annotation.class) {
bpp.setAsyncAnnotationType(customAsyncAnnotation);
}
bpp.setProxyTargetClass(this.enableAsync.get("proxyTargetClass"));
bpp.setOrder(this.enableAsync.get("order"));
return bpp;
}
}核心流程是:
@EnableAsync→@Import(AsyncConfigurationSelector.class)→ 选择ProxyAsyncConfigurationProxyAsyncConfiguration向容器注册AsyncAnnotationBeanPostProcessorAsyncAnnotationBeanPostProcessor会在 Bean 初始化后检查其是否有@Async标注的方法- 如果找到了
@Async方法,则为该 Bean 创建 AOP 代理
2.4 AsyncAnnotationBeanPostProcessor 的处理流程
AsyncAnnotationBeanPostProcessor 继承自 AbstractBeanFactoryAwareAdvisingPostProcessor,它是一个 BeanPostProcessor,在 Bean 初始化完成后检查是否需要创建异步代理。
// org.springframework.scheduling.annotation.AsyncAnnotationBeanPostProcessor
public class AsyncAnnotationBeanPostProcessor extends AbstractBeanFactoryAwareAdvisingPostProcessor {
@Override
public void setBeanFactory(BasicBeanFactory beanFactory) {
super.setBeanFactory(beanFactory);
// 创建 Advisor(切面)
AsyncAnnotationAdvisor advisor = new AsyncAnnotationAdvisor(this.executor, this.exceptionHandler);
if (this.asyncAnnotationType != null) {
advisor.setAsyncAnnotationType(this.asyncAnnotationType);
}
advisor.setBeanFactory(beanFactory);
this.advisor = advisor;
}
@Override
public Object postProcessAfterInitialization(Object bean, String beanName) {
// 如果该 Bean 需要异步代理,则创建代理
if (bean instanceof AopInfrastructureBean) {
// 基础设施 Bean 不代理
return bean;
}
return super.postProcessAfterInitialization(bean, beanName);
}
}3. AsyncConfigurer 接口
AsyncConfigurer 允许开发者 全局定制 异步执行器和异常处理器,而无需在每个 @Async 调用处重复配置。
// org.springframework.scheduling.annotation.AsyncConfigurer
public interface AsyncConfigurer {
/**
* 返回全局的异步执行器。
* 如果返回 null,Spring 会使用 SimpleAsyncTaskExecutor(简单但不适合生产)。
*/
@Nullable
default Executor getAsyncExecutor() {
return null;
}
/**
* 返回全局的异步未捕获异常处理器。
*/
@Nullable
default AsyncUncaughtExceptionHandler getAsyncUncaughtExceptionHandler() {
return null;
}
}3.1 典型实现方式
通过实现 AsyncConfigurer 可以集中管理异步线程池和异常处理逻辑:
@Configuration
@EnableAsync
public class AsyncConfig implements AsyncConfigurer {
@Override
public Executor getAsyncExecutor() {
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
executor.setCorePoolSize(10);
executor.setMaxPoolSize(50);
executor.setQueueCapacity(100);
executor.setThreadNamePrefix("async-exec-");
executor.setWaitForTasksToCompleteOnShutdown(true);
executor.setAwaitTerminationSeconds(30);
executor.initialize();
return executor;
}
@Override
public AsyncUncaughtExceptionHandler getAsyncUncaughtExceptionHandler() {
return (ex, method, params) -> {
log.error("异步方法 [{}] 执行异常: {}", method, ex.getMessage(), ex);
// 可以发送告警通知等
};
}
}3.2 更灵活的方案:使用 @Bean 定义
当需要 多个不同配置的线程池 时,可以放弃实现 AsyncConfigurer,转而通过 @Bean 定义多个 Executor,并用 @Async("beanName") 指定使用哪个:
@Configuration
@EnableAsync
public class MultiAsyncConfig {
@Bean("taskExecutor")
public ThreadPoolTaskExecutor taskExecutor() {
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
executor.setCorePoolSize(10);
executor.setMaxPoolSize(30);
executor.setQueueCapacity(200);
executor.setThreadNamePrefix("default-");
executor.initialize();
return executor;
}
@Bean("ioExecutor")
public ThreadPoolTaskExecutor ioExecutor() {
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
executor.setCorePoolSize(20);
executor.setMaxPoolSize(100);
executor.setQueueCapacity(500);
executor.setThreadNamePrefix("io-");
executor.initialize();
return executor;
}
@Bean("batchExecutor")
public ThreadPoolTaskExecutor batchExecutor() {
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
executor.setCorePoolSize(5);
executor.setMaxPoolSize(10);
executor.setQueueCapacity(1000);
executor.setThreadNamePrefix("batch-");
executor.initialize();
return executor;
}
}使用时通过注解的 value 属性指定:
@Service
public class OrderService {
@Async("ioExecutor")
public CompletableFuture<Order> queryOrderAsync(Long orderId) {
// IO 密集型任务使用 ioExecutor
}
@Async("batchExecutor")
public void batchProcess(List<Long> ids) {
// 批处理任务使用 batchExecutor
}
}4. ThreadPoolTaskExecutor 核心参数详解
ThreadPoolTaskExecutor 是对 java.util.concurrent.ThreadPoolExecutor 的 Spring 风格封装,理解其参数是合理配置异步线程池的前提。
4.1 参数一览
| 参数名 | 类型 | 默认值 | 说明 |
|---|---|---|---|
corePoolSize | int | 1 | 核心线程数,即使空闲也保留的线程数 |
maxPoolSize | int | Integer.MAX_VALUE | 最大线程数 |
queueCapacity | int | Integer.MAX_VALUE | 任务队列容量 |
keepAliveSeconds | int | 60 | 非核心线程空闲存活时间(秒) |
allowCoreThreadTimeOut | boolean | false | 是否允许核心线程超时回收 |
threadNamePrefix | String | 默认类名 | 线程名称前缀 |
rejectionPolicy | RejectedExecutionHandler | AbortPolicy | 拒绝策略 |
taskDecorator | TaskDecorator | null | 任务装饰器,用于上下文传递 |
waitForTasksToCompleteOnShutdown | boolean | false | 关闭时是否等待任务完成 |
awaitTerminationSeconds | int | 0 | 关闭时等待的最大秒数 |
4.2 线程池扩容机制
ThreadPoolTaskExecutor 的线程池扩容流程如下:
- 当一个任务提交时,如果当前运行的线程数 小于
corePoolSize,则创建新核心线程执行 - 如果当前运行的线程数 等于
corePoolSize,任务会被放入 阻塞队列 - 如果队列 已满,且当前运行的线程数 小于
maxPoolSize,则创建新非核心线程执行 - 如果队列已满且线程数达到
maxPoolSize,触发 拒绝策略
// java.util.concurrent.ThreadPoolExecutor 的任务提交流程
public void execute(Runnable command) {
int c = ctl.get();
// 1. 如果工作线程数 < corePoolSize,创建核心线程
if (workerCountOf(c) < corePoolSize) {
if (addWorker(command, true))
return;
c = ctl.get();
}
// 2. 如果线程池处于运行状态,尝试加入队列
if (isRunning(c) && workQueue.offer(command)) {
int recheck = ctl.get();
if (!isRunning(recheck) && remove(command))
reject(command);
else if (workerCountOf(recheck) == 0)
addWorker(null, false);
}
// 3. 如果队列已满,尝试创建新非核心线程
else if (!addWorker(command, false))
// 4. 如果达到 maxPoolSize,执行拒绝策略
reject(command);
}4.3 queueCapacity 与队列选择
queueCapacity 决定了阻塞队列的选择:
- queueCapacity > 0:创建
LinkedBlockingQueue(有界队列),队列满后才触发扩容 - queueCapacity <= 0:创建
SynchronousQueue(无缓冲队列),任务不排队,直接尝试创建线程
// ThreadPoolTaskExecutor.createQueue()
protected BlockingQueue<Runnable> createQueue(int queueCapacity) {
if (queueCapacity > 0) {
// 有界队列,容量由 queueCapacity 指定
return new LinkedBlockingQueue<>(queueCapacity);
} else {
// 无缓冲队列,直接转交线程
return new SynchronousQueue<>();
}
}典型场景分析:
| 场景 | 推荐配置 | 说明 |
|---|---|---|
| CPU 密集型 | core=max=CPU核数+1, queue=SynchronousQueue | 避免任务积压导致 CPU 争抢 |
| IO 密集型 | core=CPU*2, max=较大值, queue=较大有界队列 | IO 等待时让出 CPU,更多线程提升吞吐 |
| 批处理任务 | core=小, max=中, queue=大有界队列 | 控制处理速度,防止打垮下游 |
| 异步请求 | core=10, max=50, queue=200 | 通用配置,兼顾响应与资源 |
4.4 拒绝策略
ThreadPoolTaskExecutor 支持四种拒绝策略:
| 拒绝策略 | 行为 |
|---|---|
ThreadPoolExecutor.AbortPolicy(默认) | 直接抛出 RejectedExecutionException |
ThreadPoolExecutor.CallerRunsPolicy | 任务在调用者线程中直接执行(反压) |
ThreadPoolExecutor.DiscardPolicy | 静默丢弃任务 |
ThreadPoolExecutor.DiscardOldestPolicy | 丢弃队列中最旧的任务,然后重试提交 |
@Bean("asyncExecutor")
public ThreadPoolTaskExecutor asyncExecutor() {
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
executor.setCorePoolSize(10);
executor.setMaxPoolSize(30);
executor.setQueueCapacity(200);
executor.setThreadNamePrefix("async-");
// 使用 CallerRunsPolicy:让调用方线程执行,实现反压
executor.setRejectedExecutionHandler(
new ThreadPoolExecutor.CallerRunsPolicy());
executor.initialize();
return executor;
}4.5 TaskDecorator —— 任务装饰器
Spring 5.0+ 引入了 TaskDecorator,它允许在任务提交到线程池时对 Runnable 进行包装,在异步线程中 传递上下文。
@Bean("asyncExecutor")
public ThreadPoolTaskExecutor asyncExecutor() {
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
executor.setCorePoolSize(10);
executor.setMaxPoolSize(30);
executor.setQueueCapacity(200);
executor.setThreadNamePrefix("async-");
// 使用 TaskDecorator 传递上下文
executor.setTaskDecorator(runnable -> {
// 获取当前线程的上下文
RequestAttributes requestAttributes = RequestContextHolder.currentRequestAttributes();
SecurityContext securityContext = SecurityContextHolder.getContext();
return () -> {
try {
// 在异步线程中恢复上下文
RequestContextHolder.setRequestAttributes(requestAttributes);
SecurityContextHolder.setContext(securityContext);
runnable.run();
} finally {
// 清理,防止内存泄漏
RequestContextHolder.resetRequestAttributes();
SecurityContextHolder.clearContext();
}
};
});
executor.initialize();
return executor;
}5. CompletableFuture + @Async 组合使用
CompletableFuture 是 Java 8 引入的异步编程利器,结合 Spring 的 @Async 可以写出既优雅又高效的异步代码。
5.1 声明式 + 函数式编程
@Service
public class OrderQueryService {
@Async("ioExecutor")
public CompletableFuture<Order> queryFromDb(Long orderId) {
// 模拟从数据库查询订单
Order order = orderRepository.findById(orderId);
return CompletableFuture.completedFuture(order);
}
@Async("ioExecutor")
public CompletableFuture<PaymentInfo> queryPayment(Long orderId) {
// 模拟从支付系统查询支付信息
PaymentInfo payment = paymentClient.query(orderId);
return CompletableFuture.completedFuture(payment);
}
@Async("ioExecutor")
public CompletableFuture<LogisticsInfo> queryLogistics(Long orderId) {
// 模拟从物流系统查询物流信息
LogisticsInfo logistics = logisticsClient.query(orderId);
return CompletableFuture.completedFuture(logistics);
}
}5.2 组合多个异步结果
@Service
public class OrderAggregationService {
@Autowired
private OrderQueryService orderQueryService;
/**
* 并行查询订单信息、支付信息、物流信息,
* 然后合并成一个 OrderDetail 返回。
*/
public OrderDetail getOrderDetail(Long orderId) throws Exception {
CompletableFuture<Order> orderFuture = orderQueryService.queryFromDb(orderId);
CompletableFuture<PaymentInfo> paymentFuture = orderQueryService.queryPayment(orderId);
CompletableFuture<LogisticsInfo> logisticsFuture = orderQueryService.queryLogistics(orderId);
// allOf 等待所有异步任务完成
CompletableFuture<Void> allFutures = CompletableFuture.allOf(orderFuture, paymentFuture, logisticsFuture);
// 合并结果
return allFutures.thenApply(v -> {
Order order = orderFuture.join();
PaymentInfo payment = paymentFuture.join();
LogisticsInfo logistics = logisticsFuture.join();
OrderDetail detail = new OrderDetail();
detail.setOrder(order);
detail.setPayment(payment);
detail.setLogistics(logistics);
return detail;
}).get(10, TimeUnit.SECONDS); // 设置超时,防止无限等待
}
}5.3 异常处理 & 超时控制
public OrderDetail getOrderDetailWithFallback(Long orderId) {
CompletableFuture<Order> orderFuture = orderQueryService.queryFromDb(orderId)
.exceptionally(ex -> {
log.error("查询订单失败", ex);
return Order.empty(); // 降级返回空订单
})
.orTimeout(3, TimeUnit.SECONDS) // 3秒超时
.exceptionally(ex -> {
log.warn("查询订单超时", ex);
return Order.empty();
});
// 或者使用 completeOnTimeout 设置超时默认值
CompletableFuture<PaymentInfo> paymentFuture = orderQueryService.queryPayment(orderId)
.completeOnTimeout(PaymentInfo.empty(), 3, TimeUnit.SECONDS);
// 任意一个完成后就继续(取第一个返回的结果)
CompletableFuture<Order> anyFuture = CompletableFuture.anyOf(
orderQueryService.queryFromDb(orderId),
orderQueryService.queryFromCache(orderId)
).thenApply(result -> (Order) result);
return anyFuture.get();
}5.4 @Async 方法需避免的陷阱
// ❌ 错误:返回 CompletableFuture 但没有在 @Async 方法中立即返回
@Async
public CompletableFuture<Order> queryOrder(Long orderId) {
Order order = orderRepository.findById(orderId);
// 这里已经执行了异步逻辑,但返回的是已经完成的 future
return CompletableFuture.completedFuture(order);
}
// ✅ 正确:@Async 方法中做实际工作
@Async
public CompletableFuture<Order> queryOrder(Long orderId) {
// @Async AOP 拦截后,这个方法在异步线程中执行
Order order = orderRepository.findById(orderId);
return CompletableFuture.completedFuture(order);
}
// ❌ 注意:不要在一个 @Async 方法中用CompletableFuture.supplyAsync 再套一层线程池
@Async
public CompletableFuture<Order> queryOrder(Long orderId) {
return CompletableFuture.supplyAsync(() -> {
// 这里已经在异步线程中了,supplyAsync 又创建了新的异步
return orderRepository.findById(orderId);
});
}6. 异步异常处理
@Async 方法的异常处理与普通方法不同——因为异常发生在异步线程中,调用方的 try-catch 无法捕获。
6.1 有返回值的方法异常处理
对于返回 Future/CompletableFuture 的方法,异常会在调用 get()/join() 时抛出:
// 调用方处理异步方法异常
public void processOrder() {
CompletableFuture<Order> future = orderService.queryOrderAsync(123L);
try {
Order order = future.get(5, TimeUnit.SECONDS);
// 正常处理
} catch (ExecutionException e) {
// 异步方法中抛出的异常被包装在 ExecutionException 中
Throwable cause = e.getCause();
log.error("异步查询订单失败", cause);
} catch (TimeoutException e) {
log.error("异步查询订单超时", e);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
log.error("异步查询被中断", e);
}
}借助 CompletableFuture 的异常处理链式方法也能优雅处理:
@Async
public CompletableFuture<Order> queryOrderAsync(Long orderId) {
return CompletableFuture.completedFuture(orderRepository.findById(orderId));
}
// 使用 exceptionally 处理异常
orderService.queryOrderAsync(123L)
.orTimeout(5, TimeUnit.SECONDS)
.exceptionally(ex -> {
log.error("查询失败", ex);
return Order.empty();
})
.thenAccept(order -> {
if (!order.isEmpty()) {
// 正常处理
}
});6.2 无返回值的方法异常处理 —— AsyncUncaughtExceptionHandler
对于返回 void 的 @Async 方法,异步线程中的异常不会被调用方感知,必须通过 AsyncUncaughtExceptionHandler 处理。
@Configuration
@EnableAsync
public class AsyncExceptionConfig implements AsyncConfigurer {
@Override
public Executor getAsyncExecutor() {
// ...
}
@Override
public AsyncUncaughtExceptionHandler getAsyncUncaughtExceptionHandler() {
return new AsyncUncaughtExceptionHandler() {
@Override
public void handleUncaughtException(Throwable ex, Method method, Object... params) {
log.error("异步方法 [{}] 执行异常,参数: {}",
method, Arrays.toString(params), ex);
// 可以发送告警
alertService.sendAlert("异步执行异常", ex.getMessage());
}
};
}
}6.3 AsyncUncaughtExceptionHandler 源码分析
// org.springframework.aop.interceptor.AsyncUncaughtExceptionHandler
@FunctionalInterface
public interface AsyncUncaughtExceptionHandler {
/**
* 处理 @Async 方法中未捕获的异常。
*
* @param ex 抛出的异常
* @param method 被调用的异步方法
* @param params 方法参数
*/
void handleUncaughtException(Throwable ex, Method method, Object... params);
}Spring 还提供了一个简单的日志实现 SimpleAsyncUncaughtExceptionHandler,默认情况下如果用户未自定义,Spring 会使用它:
// org.springframework.aop.interceptor.SimpleAsyncUncaughtExceptionHandler
public class SimpleAsyncUncaughtExceptionHandler implements AsyncUncaughtExceptionHandler {
private static final Log logger = LogFactory.getLog(SimpleAsyncUncaughtExceptionHandler.class);
@Override
public void handleUncaughtException(Throwable ex, Method method, Object... params) {
if (logger.isErrorEnabled()) {
logger.error("异步方法 '" + method.toGenericString()
+ "' 执行时发生未捕获异常,参数: " + Arrays.toString(params), ex);
}
}
}7. 上下文传递问题
这是 Spring 异步编程中最常见的陷阱之一:线程切换导致上下文丢失。
7.1 问题复现
@Service
public class ContextAwareService {
@Async
public void asyncMethod() {
// ❌ 在主线程中设置 RequestAttributes,但在异步线程中为 null
RequestAttributes attributes = RequestContextHolder.getRequestAttributes();
// attributes == null !
// ❌ SecurityContext 同样丢失
SecurityContext context = SecurityContextHolder.getContext();
// context 是空的
}
}
@RestController
public class TestController {
@Autowired
private ContextAwareService service;
@GetMapping("/test")
public String test() {
// 主线程:RequestContextHolder 中有值
service.asyncMethod(); // 异步线程中上下文丢失
return "ok";
}
}7.2 上下文丢失的根本原因
RequestContextHolder 和 SecurityContextHolder 默认使用 ThreadLocal 存储数据,而线程池中的线程与 Web 请求线程(Tomcat 线程)不同,ThreadLocal 中的数据自然不可见。
// RequestContextHolder —— 使用 ThreadLocal 存储
public abstract class RequestContextHolder {
private static final ThreadLocal<RequestAttributes> requestAttributesHolder =
new NamedThreadLocal<>("Request attributes");
private static final NamedInheritableThreadLocal<RequestAttributes> inheritableRequestAttributesHolder =
new NamedInheritableThreadLocal<>("Request context");
public static void setRequestAttributes(@Nullable RequestAttributes attributes) {
requestAttributesHolder.set(attributes);
}
@Nullable
public static RequestAttributes getRequestAttributes() {
RequestAttributes attributes = requestAttributesHolder.get();
if (attributes == null) {
attributes = inheritableRequestAttributesHolder.get();
}
return attributes;
}
}7.3 解决方案一:TaskDecorator(推荐)
Spring 5.0+ 的 TaskDecorator 是解决上下文传递的最优雅方式:
@Configuration
@EnableAsync
public class AsyncConfig implements AsyncConfigurer {
@Override
public Executor getAsyncExecutor() {
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
executor.setCorePoolSize(10);
executor.setMaxPoolSize(30);
executor.setQueueCapacity(200);
executor.setThreadNamePrefix("async-");
// 通过 TaskDecorator 传递上下文
executor.setTaskDecorator(runnable -> {
// 捕获主线程的上下文
Map<String, String> contextMap = MDC.getCopyOfContextMap();
RequestAttributes requestAttributes = RequestContextHolder.currentRequestAttributes();
SecurityContext securityContext = SecurityContextHolder.getContext();
return () -> {
try {
// 设置到异步线程
MDC.setContextMap(contextMap);
RequestContextHolder.setRequestAttributes(requestAttributes);
SecurityContextHolder.setContext(securityContext);
runnable.run();
} finally {
// 清理,防止内存泄漏
MDC.clear();
RequestContextHolder.resetRequestAttributes();
SecurityContextHolder.clearContext();
}
};
});
executor.initialize();
return executor;
}
}7.4 解决方案二:继承性 ThreadLocal
在创建线程池时,使用 InheritableThreadLocal 可以让子线程继承父线程的上下文:
// 理论上 InheritableThreadLocal 可以让子线程继承父线程的值
// 但线程池中的线程是复用的,第二次执行时上下文仍然不正确注意
InheritableThreadLocal 仅在线程 创建时 传递一次上下文。线程池复用线程的场景下,新任务会使用旧线程残存的上下文,反而会导致数据错乱,因此不推荐此方案。
7.5 解决方案三:Spring 的 RequestContextFilter
对于 RequestAttributes,Spring Web 提供了 RequestContextFilter,它可以处理部分场景下的上下文丢失:
// 在 web.xml 或 Java Config 中配置
@Bean
public FilterRegistrationBean<RequestContextFilter> requestContextFilter() {
FilterRegistrationBean<RequestContextFilter> registration = new FilterRegistrationBean<>();
registration.setFilter(new RequestContextFilter());
registration.addUrlPatterns("/*");
registration.setOrder(Ordered.HIGHEST_PRECEDENCE);
return registration;
}但这种方式对 @Async 线程池的场景依然不够,建议始终使用 TaskDecorator 方案。
7.6 解决方案四:Hystrix / Sentinel 的上下文传播机制
当使用 Hystrix 或 Sentinel 时,它们的线程隔离机制会进一步加剧上下文丢失问题。通常它们提供了自定义的并发策略来处理:
// 示例:Hystrix 自定义并发策略
public class SpringSecurityConcurrencyStrategy extends HystrixConcurrencyStrategy {
@Override
public <T> Callable<T> wrapCallable(Callable<T> callable) {
SecurityContext securityContext = SecurityContextHolder.getContext();
RequestAttributes requestAttributes = RequestContextHolder.getRequestAttributes();
return () -> {
try {
SecurityContextHolder.setContext(securityContext);
RequestContextHolder.setRequestAttributes(requestAttributes);
return callable.call();
} finally {
SecurityContextHolder.clearContext();
RequestContextHolder.resetRequestAttributes();
}
};
}
}8. 源码分析:AsyncExecutionAspectSupport.doSubmit() 完整流程
AsyncExecutionAspectSupport 是整个 @Async 拦截执行的核心类,它位于 org.springframework.aop.interceptor 包中。
8.1 类的层次结构
AsyncExecutionAspectSupport (抽象基类)
├── AsyncExecutionInterceptor (实现 MethodInterceptor)
└── AnnotationAsyncExecutionInterceptor (添加注解解析支持)8.2 doSubmit() —— 异步提交流程的核心
doSubmit() 方法是异步执行的 核心调度器,它负责选择执行器、提交任务、处理返回值。
// org.springframework.aop.interceptor.AsyncExecutionAspectSupport
@Nullable
protected Object doSubmit(Callable<Object> task, AsyncTaskExecutor executor, Class<?> returnType) {
// ──── 分支 1: 返回类型是 CompletableFuture ────
if (CompletableFuture.class.isAssignableFrom(returnType)) {
// 使用 CompletableFuture.supplyAsync 提交到线程池
return CompletableFuture.supplyAsync(() -> {
try {
return task.call();
} catch (Throwable ex) {
throw new CompletionException(ex);
}
}, executor);
}
// ──── 分支 2: 返回类型是 ListenableFuture ────
else if (ListenableFuture.class.isAssignableFrom(returnType)) {
return ((AsyncListenableTaskExecutor) executor).submitListenable(task);
}
// ──── 分支 3: 返回类型是 Future ────
else if (Future.class.isAssignableFrom(returnType)) {
return executor.submit(task);
}
// ──── 分支 4: 返回 void 或其他类型 ────
else {
executor.submit(task);
// 返回 null —— 调用方无法获取异步结果
return null;
}
}8.3 选择执行器的过程 —— determineAsyncExecutor()
doSubmit() 被执行前,拦截器需要先确定使用哪个 Executor:
// AsyncExecutionAspectSupport.determineAsyncExecutor()
@Nullable
protected AsyncTaskExecutor determineAsyncExecutor(Method method) {
// 1. 先查缓存
AsyncTaskExecutor executor = this.executors.get(method);
if (executor == null) {
// 2. 解析 @Async 注解,获取 value(指定的执行器名称)
String qualifier = getExecutorQualifier(method);
if (StringUtils.hasLength(qualifier)) {
// 3. 按名称查找指定的 Executor Bean
executor = determineExecutorByQualifier(qualifier);
if (executor == null) {
// 指定的执行器不存在则抛出异常
throw new IllegalStateException(
"找不到名称为 [" + qualifier + "] 的 Executor Bean");
}
} else {
// 4. 没有指定名称 → 使用默认的执行器
executor = getDefaultExecutor(this.defaultExecutor, method);
}
// 5. 放入缓存
this.executors.put(method, executor);
}
return executor;
}8.4 默认执行器的选择逻辑 —— getDefaultExecutor()
// AsyncExecutionAspectSupport.getDefaultExecutor()
protected AsyncTaskExecutor getDefaultExecutor(@Nullable Executor defaultExecutor, Method method) {
// 使用 AsyncConfigurer 提供的执行器
if (defaultExecutor != null) {
return new TaskExecutorAdapter(defaultExecutor);
}
// 从 BeanFactory 中查找唯一的 Executor
if (this.beanFactory != null) {
try {
Map<String, Executor> beans = BeanFactoryUtils.beansOfTypeIncludingAncestors(
this.beanFactory, Executor.class);
if (!beans.isEmpty()) {
if (beans.size() == 1) {
// 只有一个 Executor Bean,直接使用
return new TaskExecutorAdapter(beans.values().iterator().next());
} else {
// 有多个 Executor Bean,尝试找 "taskExecutor"
Executor executor = this.beanFactory.getBean("taskExecutor", Executor.class);
return new TaskExecutorAdapter(executor);
}
}
} catch (NoUniqueBeanDefinitionException ex) {
// 多个 Executor Bean 且没有叫 "taskExecutor" 的
} catch (NoSuchBeanDefinitionException ex) {
// 没有 Executor Bean
}
}
// 以上都不满足 → 使用 SimpleAsyncTaskExecutor(生产环境不推荐!)
return new TaskExecutorAdapter(new SimpleAsyncTaskExecutor());
}8.5 完整拦截器调用链路
@Async 方法的调用
│
▼
AsyncExecutionInterceptor.invoke(MethodInvocation)
│ ┌─ 获取当前方法的 Method 对象
│ └─ AsyncExecutionAspectSupport.determineAsyncExecutor(method)
│
▼
│ ┌─ 确定执行器后,构造 Callable 任务
│ └─ AsyncExecutionAspectSupport.doSubmit(task, executor, returnType)
│
▼
├─ CompletableFuture → CompletableFuture.supplyAsync(task, executor)
├─ ListenableFuture → executor.submitListenable(task)
├─ Future → executor.submit(task)
└─ void → executor.submit(task); return null;8.6 invoke() 方法的完整代码
// org.springframework.aop.interceptor.AsyncExecutionInterceptor
@Override
@Nullable
public Object invoke(MethodInvocation invocation) throws Throwable {
Class<?> targetClass = (invocation.getThis() != null
? AopUtils.getTargetClass(invocation.getThis()) : null);
Method specificMethod = ClassUtils.getMostSpecificMethod(
invocation.getMethod(), targetClass);
final Method userDeclaredMethod = BridgeMethodResolver.findBridgedMethod(specificMethod);
// 1. 确定要使用的异步执行器
AsyncTaskExecutor executor = determineAsyncExecutor(userDeclaredMethod);
if (executor == null) {
// 如果找不到执行器,同步执行(降级)
return invocation.proceed();
}
// 2. 构造 Callable 任务
Callable<Object> task = () -> {
try {
// 实际的方法调用
Object result = invocation.proceed();
// 如果返回类型是 Future 但结果不是 Future,包装一下
if (result instanceof Future) {
return ((Future<?>) result).get();
}
return result;
} catch (ExecutionException ex) {
// 处理嵌套的 ExecutionException
handleError(ex.getCause(), userDeclaredMethod, invocation.getArguments());
return null;
} catch (Throwable ex) {
// 处理其他异常
handleError(ex, userDeclaredMethod, invocation.getArguments());
return null;
}
};
// 3. 提交到 doSubmit 核心调度方法
return doSubmit(task, executor, invocation.getMethod().getReturnType());
}
/**
* 处理异步方法中抛出的异常。
* 对于返回 Future 的方法,异常会在 get() 时抛出。
* 对于 void 方法,通过 AsyncUncaughtExceptionHandler 处理。
*/
protected void handleError(Throwable ex, Method method, Object... params) throws Exception {
if (Future.class.isAssignableFrom(method.getReturnType())) {
// Future 方法的异常会由 get() 抛出,这里不处理
throw new Exception(ex);
} else {
// void 方法的异常 → 交给 AsyncUncaughtExceptionHandler
try {
this.exceptionHandler.obtain().handleUncaughtException(ex, method, params);
} catch (Throwable ignored) {
}
}
}9. 实战案例:异步导出 10 万行 Excel + 实时进度 WebSocket 推送
本节实现一个完整的异步导出功能,包含:
- 后台线程异步生成 Excel 文件
- 通过 WebSocket 实时推送导出进度
- 支持取消、查询进度、错误处理
- 使用 CompletableFuture 管理异步任务
9.1 导出进度状态模型
// ExportProgress.java
public class ExportProgress {
/** 任务 ID */
private String taskId;
/** 当前进度(0-100) */
private int progress;
/** 已处理行数 */
private int processedRows;
/** 总行数 */
private int totalRows;
/** 状态:PENDING / PROCESSING / COMPLETED / FAILED / CANCELLED */
private String status;
/** 错误消息 */
private String errorMessage;
/** 下载地址(完成后的文件路径) */
private String downloadUrl;
// getters & setters ...
public static ExportProgress of(String taskId, int totalRows) {
ExportProgress p = new ExportProgress();
p.setTaskId(taskId);
p.setTotalRows(totalRows);
p.setProgress(0);
p.setStatus("PENDING");
return p;
}
}