问题:异步任务堆积时,调用线程为何被拖下水
在高并发的服务中,常使用 @Async 将耗时操作(如发送通知、生成报表)提交到独立线程池执行,避免阻塞 Web 请求线程。但当任务提交速度超过处理速度时,线程池会逐渐被占满,队列也会堆满。此时新提交的任务会触发拒绝策略。Spring 的 ThreadPoolTaskExecutor 默认使用 AbortPolicy,直接抛出 RejectedExecutionException,导致任务丢失。为了“不丢任务”,很多团队改用 CallerRunsPolicy,让提交任务的线程(通常是请求线程)自己执行被拒绝的任务。
这个策略看似解决了丢失问题,却引入了一个隐蔽的代价:当任务堆积严重时,调用线程会被迫执行耗时任务,阻塞在提交点。更糟的是,如果任务内部又依赖线程池中的其他任务(例如等待一个 CompletableFuture 完成),就可能形成线程互相等待的死锁。下面通过一个具体场景展开。
贯穿场景:高并发下异步任务堆积导致的服务雪崩
假设有一个订单服务,收到请求后需要异步发送通知邮件,并返回 CompletableFuture 给调用方等待结果。核心代码如下:
@Service
public class NotificationService {
@Async("taskExecutor")
public CompletableFuture<String> sendEmail(String orderId) {
// 模拟耗时操作
try { Thread.sleep(1000); } catch (InterruptedException e) { Thread.currentThread().interrupt(); }
return CompletableFuture.completedFuture("sent:" + orderId);
}
}
线程池配置如下:
@Bean("taskExecutor")
public Executor taskExecutor() {
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
executor.setCorePoolSize(2);
executor.setMaxPoolSize(2);
executor.setQueueCapacity(10);
executor.setThreadNamePrefix("taskExecutor-");
executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy());
executor.initialize();
return executor;
}
在高峰期,大量请求同时进来,每个请求都会调用 sendEmail 并立即 join() 等待结果。当线程池的两个线程都被占满,队列也堆满 10 个任务后,第 13 个任务触发 CallerRunsPolicy,由提交它的 Web 线程(Tomcat 工作线程)直接执行 sendEmail。
此时,Web 线程被阻塞在 sendEmail 的 sleep 中,无法继续接收新请求。而线程池中的两个线程可能在执行其他任务的 join() 等待,等待队列中的任务完成——但队列中的任务又需要线程池线程来执行。于是形成死锁:线程池线程等待队列任务,队列任务等待线程池线程,Web 线程则被 CallerRunsPolicy 拖入耗时任务。最终整个服务线程耗尽,表现为服务雪崩。
ThreadPoolTaskExecutor 的工作流程与队列交互
要理解上述死锁,需要先看清 ThreadPoolTaskExecutor 的任务调度顺序。它基于 ThreadPoolExecutor,处理流程如下:
- 当提交任务时,如果当前运行线程数小于
corePoolSize,直接创建新线程执行任务。 - 如果线程数已达到
corePoolSize,任务会先放入工作队列。 - 当队列满时,如果线程数小于
maximumPoolSize,创建新线程执行任务。 - 当线程数达到
maximumPoolSize且队列也满时,触发RejectedExecutionHandler。
这个顺序意味着:队列容量和线程池大小共同决定了“何时触发拒绝”。在 corePoolSize == maximumPoolSize 时,队列一旦满就会立即触发拒绝策略。
下面用流程图展示任务提交到拒绝的完整路径:
flowchart TD
A[提交任务] --> B{线程数 < corePoolSize?}
B -- 是 --> C[创建新线程执行]
B -- 否 --> D{队列未满?}
D -- 是 --> E[放入队列等待]
D -- 否 --> F{线程数 < maximumPoolSize?}
F -- 是 --> G[创建新线程执行]
F -- 否 --> H[触发拒绝策略]
H --> I[CallerRunsPolicy 由调用线程执行]
I --> J[调用线程阻塞在耗时任务]
J --> K[可能与其他线程互相等待形成死锁]
图中关键转折点是队列满且线程数达到上限,此时拒绝策略接管。CallerRunsPolicy 将任务回退给调用线程,调用线程的执行路径与线程池完全隔离,但任务本身的逻辑可能依赖线程池中的资源,这正是死锁的根源。
CallerRunsPolicy 的源码行为与上下文丢失
CallerRunsPolicy 的源码非常简洁:
public static class CallerRunsPolicy implements RejectedExecutionHandler {
public CallerRunsPolicy() { }
public void rejectedExecution(Runnable r, ThreadPoolExecutor e) {
if (!e.isShutdown()) {
r.run(); // 直接在当前线程执行任务
}
}
}
它只在线程池未关闭时执行任务,否则丢弃。这个策略的代价是:
- 调用线程阻塞:提交任务的线程(如 Tomcat 工作线程)被占用执行耗时任务,无法继续处理新请求,导致请求积压。
- 上下文丢失:任务在提交时可能携带了请求上下文(如用户信息、TraceId),但
CallerRunsPolicy直接调用r.run(),没有经过线程池的ThreadFactory或装饰器,因此无法自动传播上下文。如果任务内部依赖 ThreadLocal 中的上下文,就会获取不到。
例如,在 Spring 中,@Async 方法通常通过 TaskDecorator 来传播上下文。但 CallerRunsPolicy 绕过了 TaskDecorator,因为 r.run() 直接执行,没有经过 execute 的包装。这会导致在调用线程中执行的任务与在线程池中执行的任务上下文不一致。
死锁的根源:任务依赖与线程池占满
回到贯穿场景,死锁的形成需要两个条件:
- 任务之间存在依赖:线程池中的任务在执行
join()等待其他任务完成,而等待的任务还在队列中。 - 线程池被完全占满:没有空闲线程去执行队列中的任务。
当这两个条件同时满足时,队列中的任务永远无法获得线程,线程池线程又都在等待队列任务,形成循环等待。CallerRunsPolicy 虽然让调用线程执行了部分任务,但如果调用线程也依赖线程池中的任务,同样会陷入等待。
在示例代码中,CompletableFuture.allOf(...).join() 会阻塞当前线程直到所有子任务完成。如果子任务被提交到同一个线程池,而线程池已满,就会发生上述死锁。
解决方案:自定义拒绝策略与上下文传播
针对不同场景,可以选择不同的解决思路。
方案一:调整队列与线程池大小
最简单的做法是让队列容量小于线程池大小,这样任务更容易被新线程处理,而不是堆积在队列中。例如,将 queueCapacity 设为 1,corePoolSize 设为 2。但这并不能根治依赖型任务的死锁,只能降低概率。
方案二:自定义拒绝策略,将任务持久化或转移
如果任务不允许丢失,可以将被拒绝的任务持久化到数据库、Redis 或消息队列,由后台线程慢慢处理。例如,实现一个 RejectedExecutionHandler,在拒绝时把任务写入消息队列:
public class MqRejectedExecutionHandler implements RejectedExecutionHandler {
@Override
public void rejectedExecution(Runnable r, ThreadPoolExecutor executor) {
// 将任务序列化后发送到消息队列
// 伪代码:mqProducer.send(r);
}
}
这种方案适合对实时性要求不高的任务,但需要额外引入中间件,并处理任务序列化与幂等。
方案三:使用 TaskDecorator 传播上下文
对于上下文丢失问题,可以通过 TaskDecorator 在任务提交时包装,将 ThreadLocal 中的上下文传递到线程池线程。ThreadPoolTaskExecutor 支持设置 TaskDecorator:
executor.setTaskDecorator(runnable -> {
Map<String, String> context = MDC.getCopyOfContextMap();
return () -> {
MDC.setContextMap(context);
try {
runnable.run();
} finally {
MDC.clear();
}
};
});
但注意,CallerRunsPolicy 不会经过 TaskDecorator,因此需要在自定义拒绝策略中手动应用装饰器。
方案四:避免任务内部依赖线程池
最根本的做法是设计异步任务时,避免在任务内部 join() 等待其他任务。如果必须等待,应使用独立的线程池或使用 CompletableFuture 的异步回调,而不是阻塞等待。
拒绝策略对比与选择决策
下表对比了四种拒绝策略的适用场景和风险:
| 策略 | 行为 | 优点 | 缺点 | 适用场景 |
|---|---|---|---|---|
| AbortPolicy | 抛出异常 | 快速失败,暴露问题 | 任务丢失 | 任务可丢失或可重试 |
| CallerRunsPolicy | 调用线程执行 | 不丢任务,降低提交速度 | 调用线程阻塞,可能死锁 | 任务独立且耗时短 |
| DiscardPolicy | 直接丢弃 | 不抛异常 | 任务静默丢失 | 允许丢弃的任务 |
| DiscardOldestPolicy | 丢弃最旧任务 | 保留新任务 | 可能丢弃重要任务 | 任务有时效性 |
选择拒绝策略时,需要权衡任务的重要性、执行时间、依赖关系和系统负载。如果任务必须执行且耗时短,CallerRunsPolicy 可以接受;如果任务耗时较长或有依赖,应使用持久化方案。
生产环境中的可观测性与降级
在实施上述方案时,需要监控关键指标:
- 线程池活跃数:通过
ThreadPoolTaskExecutor的getActiveCount()监控。 - 队列大小:
getQueue().size()反映堆积情况。 - 拒绝次数:自定义拒绝策略中增加计数器。
- 线程池线程状态:使用 jstack 或 JFR 查看线程是否长时间 WAITING。
当队列持续增长或拒绝次数增加时,应触发降级:例如,暂时关闭非核心业务,或动态调整线程池参数。
总结与适用边界
CallerRunsPolicy 是一个“双刃剑”:它保证了任务不丢失,但可能将调用线程拖入耗时任务,导致阻塞甚至死锁。在任务独立且执行时间短的情况下,它是合理的选择;但在任务存在依赖或执行时间不可控时,应改用持久化或消息队列方案。
上下文传播问题可以通过 TaskDecorator 解决,但需要确保拒绝策略也应用装饰器。最终,选择哪种策略取决于任务特性、系统负载和可用资源,没有一劳永逸的方案。