线程池故障排查实战:队列堆积与拒绝策略
线程池故障排查实战:队列堆积与拒绝策略
线程池故障排查实战:队列堆积与拒绝策略是Java并发编程的核心概念,它为程序提供了并行执行的能力,是提升系统性能的关键。
本文系统介绍了线程池故障排查实战:队列堆积与拒绝策略的原理、使用方式和常见问题,帮助你掌握Java并发编程。
在排查问题之前,先快速回顾线程池的核心参数:
public ThreadPoolExecutor(
int corePoolSize, // 核心线程数
int maximumPoolSize, // 最大线程数
long keepAliveTime, // 空闲线程存活时间
TimeUnit unit, // 时间单位
BlockingQueue<Runnable> workQueue, // 工作队列
ThreadFactory threadFactory, // 线程工厂
RejectedExecutionHandler handler // 拒绝策略
)线程池的执行流程:
提交任务
│
├─ 当前线程数 < corePoolSize?──→ 创建核心线程执行
│
├─ 核心线程满了 → 任务放入 workQueue
│
├─ workQueue 满了 → 创建非核心线程(直到 maximumPoolSize)
│
└─ 线程数 = maximumPoolSize 且队列满 → 执行拒绝策略这个流程非常重要,很多线程池问题都源于对执行流程的误解。
二、故障场景一:队列堆积
2.1 问题现象
监控系统告警:
告警: ThreadPool "orderProcessor" queue size = 500000
告警: ThreadPool "orderProcessor" active threads = 8 (max=8)
告警: 任务处理延迟 = 120秒(正常 < 1秒)用户反馈:下单后状态更新慢了 2 分钟。
2.2 排查步骤
第一步:查看线程池状态
// 方案一:通过 JMX 暴露线程池指标
@Bean
public MeterBinder threadPoolMetrics(ThreadPoolExecutor orderExecutor) {
return registry -> {
Gauge.builder("threadpool.active.count", orderExecutor, ThreadPoolExecutor::getActiveCount)
.tag("name", "orderProcessor")
.register(registry);
Gauge.builder("threadpool.queue.size", orderExecutor, e -> e.getQueue().size())
.tag("name", "orderProcessor")
.register(registry);
Gauge.builder("threadpool.pool.size", orderExecutor, ThreadPoolExecutor::getPoolSize)
.tag("name", "orderProcessor")
.register(registry);
};
}第二步:用 jstack 查看线程在干什么
jstack <PID> | grep -A 10 "orderProcessor"输出:
"orderProcessor-1-thread-3" #45 prio=5 os_prio=0
java.lang.Thread.State: RUNNABLE
at java.net.SocketInputStream.socketRead0(Native Method)
at com.example.service.OrderService.saveOrder(OrderService.java:87)
← 线程在等数据库响应!所有 8 个核心线程都在等数据库 IO,任务积压在队列中。
第三步:查看数据库状态
-- 查看当前连接
SHOW PROCESSLIST;
-- 发现大量慢查询
SELECT * FROM information_schema.PROCESSLIST WHERE TIME > 5;数据库慢查询导致每个任务处理时间从 50ms 飙升到 5 秒,8 个线程的处理速度跟不上任务提交速度,队列开始堆积。
2.3 根因分析
正常情况:
任务提交速率: 100/秒
处理能力: 8线程 × 50ms/任务 = 160/秒 → 稳定
异常情况:
任务提交速率: 100/秒
处理能力: 8线程 × 5000ms/任务 = 1.6/秒 → 堆积!
队列增长: 100 - 1.6 = 98.4/秒
1 分钟后队列: 98.4 × 60 = 5904
10 分钟后队列: 59040 ← 爆了2.4 修复方案
立即止血:扩容线程 + 限制队列
// 原配置
ThreadPoolExecutor executor = new ThreadPoolExecutor(
8, // corePoolSize
16, // maximumPoolSize
60, TimeUnit.SECONDS,
new LinkedBlockingQueue<>(100000), // 大队列
new ThreadPoolExecutor.CallerRunsPolicy()
);
// 临时调整
executor.setCorePoolSize(32); // 扩大核心线程
executor.setMaximumPoolSize(64); // 扩大最大线程注意:如果根因是数据库慢,扩线程没用,反而会让数据库更慢。止血应该是限流。
限流降级:
// 限制提交速率
RateLimiter rateLimiter = RateLimiter.create(50); // 50 任务/秒
public void submitOrder(Order order) {
if (!rateLimiter.tryAcquire(1, TimeUnit.SECONDS)) {
throw new ServiceUnavailableException("系统繁忙,请稍后重试");
}
executor.submit(() -> processOrder(order));
}根因修复:优化数据库
// 优化前:每个订单单独查询
void processOrder(Order order) {
User user = userMapper.selectById(order.getUserId()); // 查询1
Product product = productMapper.selectById(order.getProductId()); // 查询2
Coupon coupon = couponMapper.selectById(order.getCouponId()); // 查询3
orderMapper.update(order); // 查询4
}
// 优化后:批量查询 + 减少查询次数
void processOrder(Order order) {
// 用 JOIN 一次查出所有信息
OrderDetail detail = orderMapper.selectOrderDetail(order.getId());
orderMapper.update(order);
}三、故障场景二:拒绝策略触发
3.1 四种内置拒绝策略
| 策略 | 行为 | 适用场景 |
|---|---|---|
| AbortPolicy | 抛 RejectedExecutionException | 默认,重要任务 |
| CallerRunsPolicy | 由提交线程执行 | 不想丢任务,允许降速 |
| DiscardPolicy | 静默丢弃 | 可丢弃的任务 |
| DiscardOldestPolicy | 丢弃队列最老的任务 | 只关心最新任务 |
3.2 问题场景
// 线程池配置
ThreadPoolExecutor executor = new ThreadPoolExecutor(
4, 8, 60, TimeUnit.SECONDS,
new ArrayBlockingQueue<>(100),
new ThreadPoolExecutor.AbortPolicy() // 默认拒绝策略
);
// 高并发时抛出大量异常
executor.submit(() -> processOrder(order));
// 抛出: java.util.concurrent.RejectedExecutionException用户看到的:下单失败,提示"系统繁忙"。
3.3 排查思路
// 增加线程池监控日志
public class MonitoredThreadPoolExecutor extends ThreadPoolExecutor {
private final AtomicInteger rejectedCount = new AtomicInteger(0);
@Override
protected void afterExecute(Runnable r, Throwable t) {
super.afterExecute(r, t);
log.info("线程池状态: active={}, poolSize={}, queueSize={}, completed={}",
getActiveCount(), getPoolSize(), getQueue().size(), getCompletedTaskCount());
}
@Override
public void rejectedExecution(Runnable r, ThreadPoolExecutor executor) {
rejectedCount.incrementAndGet();
log.warn("任务被拒绝!当前队列大小: {}, 拒绝次数: {}",
executor.getQueue().size(), rejectedCount.get());
super.rejectedExecution(r, executor); // 执行默认拒绝策略
}
}3.4 拒绝策略选择建议
场景一:核心业务(不能丢)
// 使用 CallerRunsPolicy,让调用线程自己执行
// 好处:自然限流,不会丢任务
// 坏处:调用线程被阻塞,可能导致接口超时
new ThreadPoolExecutor(
8, 16, 60, TimeUnit.SECONDS,
new ArrayBlockingQueue<>(200),
new ThreadPoolExecutor.CallerRunsPolicy()
);场景二:可丢弃的非核心业务(如日志上报)
// 自定义拒绝策略:记录日志后丢弃
new ThreadPoolExecutor(
2, 4, 30, TimeUnit.SECONDS,
new ArrayBlockingQueue<>(100),
(r, executor) -> {
log.warn("日志任务被丢弃,当前队列大小: {}", executor.getQueue().size());
// 不执行 r,直接丢弃
}
);场景三:需要持久化的任务
// 拒绝时持久化到数据库或 MQ,后续恢复
new ThreadPoolExecutor(
8, 16, 60, TimeUnit.SECONDS,
new ArrayBlockingQueue<>(200),
(r, executor) -> {
// 将任务序列化存入 Redis/MQ,等待后续处理
if (r instanceof OrderTask) {
OrderTask task = (OrderTask) r;
redisTemplate.opsForList().leftPush("pending:orders", JSON.toJSONString(task));
log.warn("任务溢出,已持久化到 Redis: {}", task.getOrderId());
}
}
);四、故障场景三:线程泄漏
4.1 问题现象
监控显示线程数持续增长,从初始的 8 个涨到了 500 多个,但任务量并没有增加。
4.2 排查步骤
# 查看线程数
jstack <PID> | grep "java.lang.Thread.State" | wc -l
# 输出: 532
# 查看线程名称分布
jstack <PID> | grep "^\"" | awk '{print $1}' | sort | uniq -c | sort -rn | head -20输出:
500 "pool-3-thread-*"
12 "http-nio-8080-exec-*"
8 "orderProcessor-*"
4 "GC task thread"pool-3-thread-* 有 500 个线程,这是哪里来的?
4.3 根因:每次请求都创建线程池
// 问题代码:在方法内部创建线程池
public void batchProcess(List<Item> items) {
// 每次调用都创建一个新线程池!
ExecutorService executor = Executors.newFixedThreadPool(10);
for (Item item : items) {
executor.submit(() -> processItem(item));
}
// 忘记 shutdown!线程池不会被 GC,线程一直存活
// executor.shutdown(); ← 漏了这行
}Executors.newFixedThreadPool 创建的线程默认以非守护线程运行,即使方法返回了,线程仍然存活,不会被 GC 回收。
4.4 修复方案
方案一:正确关闭线程池
public void batchProcess(List<Item> items) {
ExecutorService executor = Executors.newFixedThreadPool(10);
try {
List<Future<?>> futures = new ArrayList<>();
for (Item item : items) {
futures.add(executor.submit(() -> processItem(item)));
}
// 等待所有任务完成
for (Future<?> f : futures) {
f.get(30, TimeUnit.SECONDS);
}
} catch (Exception e) {
log.error("批量处理失败", e);
} finally {
executor.shutdown(); // 关键:关闭线程池
try {
if (!executor.awaitTermination(60, TimeUnit.SECONDS)) {
executor.shutdownNow();
}
} catch (InterruptedException e) {
executor.shutdownNow();
Thread.currentThread().interrupt();
}
}
}方案二:使用全局线程池(推荐)
// 全局线程池,Spring Bean 管理
@Configuration
public class ThreadPoolConfig {
@Bean("batchProcessExecutor")
public ThreadPoolExecutor batchProcessExecutor() {
return new ThreadPoolExecutor(
10, 20, 60, TimeUnit.SECONDS,
new LinkedBlockingQueue<>(500),
new ThreadFactoryBuilder().setNameFormat("batch-process-%d").build(),
new ThreadPoolExecutor.CallerRunsPolicy()
);
}
}
@Service
public class BatchService {
@Resource(name = "batchProcessExecutor")
private ThreadPoolExecutor executor;
public void batchProcess(List<Item> items) {
// 复用全局线程池,不需要创建/销毁
for (Item item : items) {
executor.submit(() -> processItem(item));
}
}
}五、故障场景四:线程池配置不当
5.1 CPU 密集型 vs IO 密集型
// CPU 密集型任务(计算、加解密、序列化)
// 线程数 = CPU 核心数 + 1
int cpuThreads = Runtime.getRuntime().availableProcessors() + 1;
// IO 密集型任务(数据库、HTTP 调用、文件读写)
// 线程数 = CPU 核心数 × (1 + IO等待时间/CPU时间)
// 简化公式:CPU 核心数 × 2 或更高
int ioThreads = Runtime.getRuntime().availableProcessors() * 2;更精确的公式(Little's Law):
最优线程数 = (任务总时间 / CPU时间) × CPU核心数
示例:
CPU 时间: 20ms
IO 等待: 80ms
总时间: 100ms
CPU 核心数: 8
最优线程数 = (100/20) × 8 = 405.2 常见配置错误
错误一:用 Executors 工具类创建线程池
// 危险!无界队列,可能 OOM
ExecutorService executor1 = Executors.newFixedThreadPool(10);
// 内部是 new LinkedBlockingQueue<>(Integer.MAX_VALUE)
// 危险!无界线程数,可能创建大量线程
ExecutorService executor2 = Executors.newCachedThreadPool();
// 内部是 new SynchronousQueue<>(), maximumPoolSize = Integer.MAX_VALUE阿里规范明确禁止使用 Executors 创建线程池,必须手动创建。
正确做法:
ThreadPoolExecutor executor = new ThreadPoolExecutor(
corePoolSize,
maxPoolSize,
60, TimeUnit.SECONDS,
new ArrayBlockingQueue<>(1000), // 有界队列
new ThreadFactoryBuilder()
.setNameFormat("biz-pool-%d")
.setUncaughtExceptionHandler((t, e) ->
log.error("线程 {} 未捕获异常", t.getName(), e))
.build(),
new ThreadPoolExecutor.CallerRunsPolicy()
);六、线程池监控体系
6.1 关键指标
| 指标 | 含义 | 告警阈值 |
|---|---|---|
| activeCount | 活跃线程数 | > maxPoolSize × 0.8 |
| queueSize | 队列大小 | > capacity × 0.8 |
| queueRemaining | 队列剩余容量 | < capacity × 0.2 |
| completedTaskCount | 已完成任务数 | 监控趋势 |
| rejectedCount | 拒绝任务数 | > 0 |
| poolSize | 当前线程数 | - |
| largestPoolSize | 历史最大线程数 | - |
6.2 Spring Boot 动态监控
@Component
@Slf4j
public class ThreadPoolMonitor implements ApplicationListener<ApplicationReadyEvent> {
@Resource(name = "orderExecutor")
private ThreadPoolExecutor orderExecutor;
@Override
public void onApplicationEvent(ApplicationReadyEvent event) {
ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor();
scheduler.scheduleAtFixedRate(this::logStatus, 10, 10, TimeUnit.SECONDS);
}
private void logStatus() {
log.info("线程池[orderExecutor] 状态: active={}, pool={}, queue={}, completed={}, rejected={}",
orderExecutor.getActiveCount(),
orderExecutor.getPoolSize(),
orderExecutor.getQueue().size(),
orderExecutor.getCompletedTaskCount(),
// 如果是自定义拒绝策略可以记录拒绝次数
"N/A"
);
}
}6.3 接入 Prometheus + Grafana
@Configuration
public class ThreadPoolMetricsConfig {
@Bean
public MeterBinder orderExecutorMetrics(@Qualifier("orderExecutor") ThreadPoolExecutor executor) {
return new MeterBinder() {
@Override
public void bindTo(MeterRegistry registry) {
Gauge.builder("tp.active", executor, ThreadPoolExecutor::getActiveCount)
.tag("pool", "order")
.register(registry);
Gauge.builder("tp.queue.size", executor, e -> e.getQueue().size())
.tag("pool", "order")
.register(registry);
Gauge.builder("tp.queue.capacity", executor, e -> e.getQueue().remainingCapacity() + e.getQueue().size())
.tag("pool", "order")
.register(registry);
Gauge.builder("tp.pool.size", executor, ThreadPoolExecutor::getPoolSize)
.tag("pool", "order")
.register(registry);
}
};
}
}七、线程池参数动态调整
线上修改线程池参数,不想重启应用?ThreadPoolExecutor 提供了动态修改方法:
// 动态修改核心线程数
executor.setCorePoolSize(newCore);
// 动态修改最大线程数
executor.setMaximumPoolSize(newMax);
// 设置核心线程是否允许超时回收
executor.allowCoreThreadTimeOut(true);配合配置中心(如 Nacos、Apollo)实现动态调整:
@NacosConfigListener(dataId = "thread-pool-config")
public void onConfigChange(String config) {
JSONObject json = JSON.parseObject(config);
int core = json.getIntValue("orderExecutor.core");
int max = json.getIntValue("orderExecutor.max");
log.info("动态调整线程池: core {} → {}, max {} → {}",
orderExecutor.getCorePoolSize(), core,
orderExecutor.getMaximumPoolSize(), max);
orderExecutor.setCorePoolSize(core);
orderExecutor.setMaximumPoolSize(max);
}八、面试要点总结
Q1:线程池的工作原理?
提交任务后:
- 当前线程数 < corePoolSize → 创建核心线程执行
- 核心线程满了 → 任务入队列
- 队列满了 → 创建非核心线程(到 maximumPoolSize)
- 线程和队列都满了 → 执行拒绝策略
Q2:四种拒绝策略是什么?
- AbortPolicy:抛异常(默认)
- CallerRunsPolicy:调用线程执行
- DiscardPolicy:静默丢弃
- DiscardOldestPolicy:丢弃最老的任务
Q3:为什么不推荐用 Executors 创建线程池?
newFixedThreadPool和newSingleThreadExecutor:使用无界队列(LinkedBlockingQueue),可能 OOMnewCachedThreadPool和newScheduledThreadPool:最大线程数为Integer.MAX_VALUE,可能创建大量线程
Q4:线程池线程数怎么设置?
- CPU 密集型:
CPU 核心数 + 1 - IO 密集型:
CPU 核心数 × (1 + IO时间/CPU时间) - 实际以压测结果为准
Q5:线程池队列堆积怎么排查?
- 查看队列大小和活跃线程数
- jstack 查看线程在干什么(通常在等 IO)
- 检查下游服务(数据库、RPC)是否变慢
- 检查任务处理逻辑是否有性能问题
- 临时方案:扩容线程或限流
九、总结
线程池故障排查的核心思路:
- 监控先行:线程池的 active count、queue size、rejected count 必须有监控
- 看线程在干什么:jstack 是利器,线程卡在哪里一目了然
- 区分根因和表象:队列堆积是表象,根因通常是下游变慢或任务处理慢
- 合理配置:有界队列 + 合适的拒绝策略 + 动态调整能力
- 不要用 Executors:手动创建 ThreadPoolExecutor,自己控制每个参数
记住一个经验:线程池的问题 90% 不是线程池配置的问题,而是下游服务或任务逻辑的问题。线程池只是"受害 者",真正的凶手在别处。