线程池源码级生命周期与线上排查
很多人学线程池只停在“七个参数”和“核心线程、队列、最大线程、拒绝策略”。这能应付一部分面试,但到了生产环境还不够。真正线上排查时,你会遇到:
- 明明线程池没满,接口还是慢。
- 队列一直涨,不知道是线程少还是下游慢。
submit的任务失败了,日志里却看不到异常。shutdown后服务迟迟停不掉。jstack里线程都是WAITING、TIMED_WAITING、BLOCKED,不知道分别代表什么。- 线程池越调越大,吞吐没上去,数据库反而被打满。
这篇把线程池从源码生命周期、任务流转、线程状态、监控指标、排查流程串成一条线。
学习目标
学完这一页,你要能回答:
ThreadPoolExecutor内部为什么同时管理“线程池状态”和“线程数量”。execute()、addWorker()、runWorker()、getTask()大致做什么。- 一个任务从提交、入队、执行、异常、完成到线程复用的完整过程。
- 线程为什么会被回收,为什么有的线程一直不退出。
execute和submit的异常为什么表现不同。jstack中不同线程状态如何对应线程池问题。- 队列堆积时如何判断是线程池配置问题、下游瓶颈、锁等待、CPU 瓶颈还是提交速度过快。
- 线程池参数调整为什么必须联动数据库连接池、HTTP 连接池、Redis、MQ 和下游接口。
线程池不是一个简单队列
ThreadPoolExecutor 至少同时管理五件事:
flowchart TD
A["ThreadPoolExecutor"] --> B["线程池状态<br/>RUNNING / SHUTDOWN / STOP"]
A --> C["工作线程数量<br/>workerCount"]
A --> D["任务队列<br/>workQueue"]
A --> E["Worker 集合<br/>workers"]
A --> F["拒绝策略<br/>handler"]
A --> G["线程工厂<br/>threadFactory"]如果只把线程池理解成“几个线程加一个队列”,会漏掉两个关键点:
- 线程池要判断自己是否还能接收新任务。
- 线程池要控制工作线程数量不能超过边界。
这就是为什么源码里会把运行状态和线程数量一起管理。不同 JDK 实现细节可能有差异,但思想是一致的:状态决定能不能接任务,线程数量决定能不能创建 Worker。
运行状态为什么重要
线程池常见状态:
| 状态 | 是否接收新任务 | 是否处理队列任务 | 典型来源 |
|---|---|---|---|
RUNNING | 是 | 是 | 正常运行 |
SHUTDOWN | 否 | 是 | 调用 shutdown() |
STOP | 否 | 否,并尝试中断执行中任务 | 调用 shutdownNow() |
TIDYING | 否 | 否 | 所有任务结束,准备终止 |
TERMINATED | 否 | 否 | 完全终止 |
流程图:
flowchart TD
A["RUNNING<br/>接收新任务并处理队列"] --> B["SHUTDOWN<br/>不接新任务但处理队列"]
A --> C["STOP<br/>不接新任务并中断任务"]
B --> D["TIDYING<br/>任务全部结束"]
C --> D
D --> E["TERMINATED<br/>线程池终止"]为什么需要这些状态?
| 如果没有状态控制 | 后果 |
|---|---|
| 关闭后仍然接收任务 | 应用停机时任务越积越多 |
shutdown 直接杀任务 | 已提交的订单、消息、文件处理可能丢失 |
shutdownNow 不尝试中断 | 紧急停机时长时间卡住 |
| 不区分队列任务和新任务 | 无法做到优雅停机 |
一个任务提交后的完整生命周期
以 execute() 为例,可以把流程拆成三段。
第一段:是否直接创建核心 Worker
flowchart TD
A["execute(command)"] --> B{"workerCount < corePoolSize"}
B -- "是" --> C["addWorker(command, true)"]
C --> D{"创建成功"}
D -- "是" --> E["新 Worker 执行 command"]
D -- "否" --> F["进入后续入队流程"]
B -- "否" --> F这里的重点是:只要当前工作线程数小于核心线程数,线程池倾向于创建新 Worker,而不是先排队。
第二段:核心线程满后尝试入队
flowchart TD
A["核心线程已满"] --> B{"线程池仍是 RUNNING"}
B -- "否" --> C["拒绝任务"]
B -- "是" --> D{"workQueue.offer(command)"}
D -- "成功" --> E["任务进入队列"]
E --> F["再次检查线程池状态"]
F --> G{"如果关闭则移除并拒绝"}
F --> H{"如果没有 Worker 则补一个空 Worker"}
D -- "失败" --> I["尝试创建非核心 Worker"]为什么入队后还要再次检查状态?
因为并发环境下,任务刚入队,另一个线程可能调用了 shutdown()。如果不二次检查,就可能出现线程池关闭后仍然留下新任务。
第三段:队列满后创建非核心 Worker 或拒绝
flowchart TD
A["队列已满"] --> B{"workerCount < maximumPoolSize"}
B -- "是" --> C["addWorker(command, false)"]
C --> D{"创建成功"}
D -- "是" --> E["非核心 Worker 执行任务"]
D -- "否" --> F["执行拒绝策略"]
B -- "否" --> F这就是“先核心线程、再队列、再最大线程、最后拒绝”的完整原因。
Worker 到底是什么
Worker 可以理解成线程池内部的工作单元。它不只是一个 Thread,还要记录第一个任务、执行状态,并参与线程池生命周期管理。
flowchart TD
A["Worker"] --> B["thread<br/>真正运行的线程"]
A --> C["firstTask<br/>创建 Worker 时带入的第一个任务"]
A --> D["completedTasks<br/>完成任务数"]
A --> E["锁控制<br/>避免运行中被错误中断"]为什么 Worker 需要 firstTask?
因为创建线程的目的通常就是为了立刻执行当前提交的任务。如果每次都先创建空线程,再让线程去队列里取任务,会多一步入队和竞争。firstTask 可以让新 Worker 直接执行当前任务。
runWorker 循环:线程为什么能复用
线程池复用线程的核心在 runWorker 思想:一个 Worker 执行完第一个任务后,并不马上退出,而是继续从队列取任务。
flowchart TD
A["Worker 线程启动"] --> B["取 firstTask"]
B --> C{"task 是否为空"}
C -- "否" --> D["执行 task.run()"]
D --> E["执行 afterExecute"]
E --> F["task 置空"]
F --> G["getTask 从队列取下一个任务"]
G --> C
C -- "是" --> H["Worker 退出"]
H --> I["processWorkerExit"]这就是为什么线程池不是“一个任务一个线程”。正确理解是:
一个 Worker 线程会执行多个任务,任务之间通过队列衔接。
如果任务里使用了 ThreadLocal 却不清理,下一个任务复用同一个线程时就可能读到旧值。
getTask:线程为什么会等待或退出
getTask() 的职责是从队列取任务,并决定当前 Worker 是否应该退出。
flowchart TD
A["getTask"] --> B{"线程池状态是否允许取队列任务"}
B -- "不允许" --> C["减少 workerCount 并退出"]
B -- "允许" --> D{"当前线程是否需要超时回收"}
D -- "是" --> E["poll 等待 keepAliveTime"]
D -- "否" --> F["take 一直等待任务"]
E --> G{"是否拿到任务"}
F --> G
G -- "拿到" --> H["返回任务"]
G -- "超时未拿到" --> I["判断是否退出 Worker"]这里能解释几个现象:
| 现象 | 原因 |
|---|---|
核心线程长期 WAITING | 默认核心线程不超时,队列为空时一直 take() 等任务 |
| 非核心线程一段时间后消失 | 超过 keepAliveTime 没任务,被回收 |
开启 allowCoreThreadTimeOut(true) 后核心线程也会退出 | 核心线程也使用超时等待 |
shutdown() 后线程池能处理完队列再退出 | 状态不接新任务,但允许继续取队列任务 |
execute 和 submit 异常为什么不同
execute
execute() 提交的是 Runnable。任务抛出运行时异常时,异常会从工作线程里抛出,通常可以被线程的 UncaughtExceptionHandler 或日志看到,线程池会补充新的 Worker。
pool.execute(() -> {
throw new RuntimeException("execute task failed");
});submit
submit() 会把任务包装成 FutureTask。任务异常不会直接抛到工作线程外,而是保存到 FutureTask 内部,调用 future.get() 时才抛出 ExecutionException。
Future<?> future = pool.submit(() -> {
throw new RuntimeException("submit task failed");
});
future.get(); // 这里才能感知异常对比:
| 提交方式 | 异常去哪了 | 风险 |
|---|---|---|
execute | 从任务线程抛出 | 如果没设置线程异常处理器,日志可能不清晰 |
submit | 保存到 Future | 不调用 get() 就像“没失败” |
商业项目里,异步入库、发送 MQ、生成报表这类任务如果用 submit 后不检查 Future,可能出现“任务失败但业务以为成功提交”的问题。
beforeExecute 和 afterExecute 有什么用
ThreadPoolExecutor 提供了钩子方法:
protected void beforeExecute(Thread t, Runnable r) { }
protected void afterExecute(Runnable r, Throwable t) { }
protected void terminated() { }它们适合做:
- 统计任务耗时。
- 捕获和记录任务异常。
- 清理
ThreadLocal。 - 上报线程池指标。
- 在线程池终止时释放资源。
Demo:
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
public class TraceableThreadPool extends ThreadPoolExecutor {
private static final ThreadLocal<Long> START_TIME = new ThreadLocal<Long>();
public TraceableThreadPool() {
super(
4,
8,
60,
TimeUnit.SECONDS,
new ArrayBlockingQueue<Runnable>(200),
r -> new Thread(r, "trace-worker-" + System.nanoTime()),
new CallerRunsPolicy()
);
}
@Override
protected void beforeExecute(Thread thread, Runnable task) {
START_TIME.set(System.currentTimeMillis());
}
@Override
protected void afterExecute(Runnable task, Throwable throwable) {
try {
Long start = START_TIME.get();
if (start != null) {
long cost = System.currentTimeMillis() - start;
System.out.println("task cost=" + cost + "ms");
}
if (throwable != null) {
System.out.println("task error=" + throwable.getMessage());
}
} finally {
START_TIME.remove();
}
}
}注意:如果任务是 submit() 包装出来的 FutureTask,异常可能在 Future 里,afterExecute 的 throwable 不一定直接拿到。生产中要额外判断 Runnable 是否是 Future,再调用 get() 获取异常,但不要在普通未完成任务上阻塞等待。
线程状态和线程池问题怎么对应
用 jstack 看线程池时,不要只看线程数量,要看线程都卡在哪里。
| 线程状态 | 常见含义 | 线程池场景 |
|---|---|---|
RUNNABLE | 正在运行或等待 CPU/系统调用 | CPU 计算、网络读写、本地方法 |
WAITING | 无限期等待 | 空闲 Worker 在队列 take() 等任务 |
TIMED_WAITING | 限时等待 | sleep、带超时的 poll、HTTP read timeout |
BLOCKED | 等待进入 synchronized 锁 | 多线程竞争同一把对象锁 |
WAITING on condition | 等条件满足 | 等数据库连接、队列、Future 结果 |
场景一:大量线程 WAITING 在队列
可能是正常空闲:
java.lang.Thread.State: WAITING
at java.util.concurrent.locks.LockSupport.park
at java.util.concurrent.LinkedBlockingQueue.take如果线程池没有任务,这是正常的。不要看到 WAITING 就以为线程死了。
场景二:大量线程卡在 socketRead
java.lang.Thread.State: RUNNABLE
at java.net.SocketInputStream.socketRead0通常表示线程在等网络响应,例如 HTTP、数据库、Redis 或第三方接口。此时盲目增加线程池可能把下游打得更慢。
排查方向:
- 下游接口 P95/P99 是否升高。
- 是否设置连接超时和读取超时。
- HTTP 连接池是否打满。
- 是否有重试放大流量。
场景三:大量线程等待数据库连接
常见表现:
waiting on condition
at com.zaxxer.hikari.pool.HikariPool.getConnection这说明线程池线程已经拿到执行机会,但卡在连接池。线程池调大不能解决数据库连接不足,可能造成更多线程排队。
处理方向:
- 查 SQL 是否慢。
- 查数据库连接池活跃数、等待数。
- 查事务是否过大或连接未关闭。
- 评估是否需要限流、拆批或优化索引。
场景四:大量线程 BLOCKED
说明线程在抢同一把 synchronized 锁。
处理方向:
- 找到锁对象和持有锁线程。
- 看锁内是否做了 IO、数据库、远程调用。
- 缩小锁范围。
- 用并发容器或更细粒度锁替代大锁。
线程池堆积排查决策树
flowchart TD
A["线程池任务堆积"] --> B["看提交 TPS 和完成 TPS"]
B --> C{"提交 TPS 是否长期大于完成 TPS"}
C -- "否" --> D["可能是瞬时峰值或监控窗口问题"]
C -- "是" --> E["看 activeCount"]
E --> F{"activeCount 是否接近 maximumPoolSize"}
F -- "否" --> G["查任务是否没提交、调度是否异常、核心线程是否过小"]
F -- "是" --> H["看线程栈"]
H --> I{"线程主要卡在哪里"}
I -- "CPU 计算" --> J["CPU 瓶颈,减少线程或优化算法"]
I -- "数据库连接/慢 SQL" --> K["优化 SQL、连接池、事务、限流"]
I -- "HTTP/Redis/下游接口" --> L["设置超时、熔断、隔离、降并发"]
I -- "锁等待" --> M["缩小锁范围或改并发结构"]
I -- "Future.get 等子任务" --> N["排查线程饥饿死锁"]
K --> O["调整线程池和下游容量"]
L --> O
M --> O
N --> O关键判断:
| 指标 | 说明 |
|---|---|
| 提交 TPS | 每秒进入线程池的任务数 |
| 完成 TPS | 每秒真正完成的任务数 |
| 队列长度 | 积压任务数 |
| activeCount | 正在执行任务的线程数 |
| 拒绝次数 | 线程池过载次数 |
| 单任务耗时 | 平均、P95、P99 耗时 |
| 下游耗时 | DB、Redis、HTTP、MQ 等耗时 |
如果提交 TPS 长期大于完成 TPS,队列必然越来越大。此时不是“调一个神奇参数”能解决,而是要提升完成 TPS、降低提交 TPS,或让系统明确拒绝/反压。
线程池参数调整的正确顺序
很多人一看到堆积就加 maximumPoolSize。这很危险。
正确顺序:
flowchart TD
A["线程池有问题"] --> B["先确认任务类型"]
B --> C["看线程栈确认卡点"]
C --> D["看下游容量"]
D --> E["看队列等待时间"]
E --> F["决定加线程、降并发、调队列或拒绝"]
F --> G["压测验证"]
G --> H["上线监控和回滚预案"]| 现象 | 不建议 | 建议 |
|---|---|---|
| DB 慢导致线程堆积 | 直接加线程 | 优化 SQL、限流、调连接池、拆批 |
| HTTP 下游慢 | 加大队列无限等 | 超时、熔断、隔离、降级 |
| CPU 100% | 增加线程 | 减少线程、优化算法、扩容机器 |
| 队列过大用户超时 | 继续加队列 | 缩小队列、快速失败或反压 |
| 任务有依赖等待 | 加大最大线程数 | 拆线程池或改异步编排 |
商业场景:医疗采集任务堆积
假设采集平台每分钟拉取医院接口数据,任务进入采集线程池。
flowchart TD
A["调度中心触发采集"] --> B["提交到采集线程池"]
B --> C["调用医院接口"]
C --> D["清洗和校验"]
D --> E["写 MySQL"]
E --> F["发送 MQ 同步 ES"]
F --> G["更新采集任务状态"]如果线程池队列上涨,可能原因不是线程池本身:
| 卡点 | 表现 | 处理 |
|---|---|---|
| 医院接口慢 | 线程卡在 socketRead | 设置超时、降并发、失败补偿 |
| 数据库慢 | 线程等待连接或慢 SQL | 优化索引、批量写入、缩事务 |
| ES/MQ 慢 | 发送耗时升高 | 异步化、失败表、重试和死信 |
| 清洗 CPU 高 | CPU 接近 100% | 控制线程数、优化规则、扩容 |
| 子任务等待 | Future.get 卡住 | 拆池或 CompletableFuture 编排 |
可运行简化 Demo:
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
public class CollectPoolDemo {
private static final ThreadPoolExecutor COLLECT_POOL =
new ThreadPoolExecutor(
8,
16,
60,
TimeUnit.SECONDS,
new ArrayBlockingQueue<Runnable>(500),
r -> new Thread(r, "collect-worker-" + System.nanoTime()),
new ThreadPoolExecutor.CallerRunsPolicy()
);
public static void main(String[] args) {
for (int i = 0; i < 100; i++) {
int taskId = i;
COLLECT_POOL.execute(() -> collect(taskId));
}
COLLECT_POOL.shutdown();
}
private static void collect(int taskId) {
try {
mockHospitalApi();
mockSaveDb();
System.out.println(Thread.currentThread().getName() + " finish task " + taskId);
} catch (Exception e) {
System.out.println("collect failed taskId=" + taskId + ", error=" + e.getMessage());
}
}
private static void mockHospitalApi() throws InterruptedException {
TimeUnit.MILLISECONDS.sleep(100);
}
private static void mockSaveDb() throws InterruptedException {
TimeUnit.MILLISECONDS.sleep(50);
}
}这个 Demo 不是为了模拟真实医院接口,而是体现生产原则:
- 线程池有界。
- 线程有业务名称。
- 满载时
CallerRunsPolicy反压提交方。 - 任务内部要捕获并记录失败。
- 真实项目还要加超时、幂等、失败表和告警。
面试标准回答
线程池任务堆积怎么排查
标准回答:
线程池堆积不能只看线程数。先看提交 TPS、完成 TPS、队列长度、activeCount、拒绝次数和单任务耗时。如果提交速度长期大于完成速度,队列一定上涨。然后用 jstack 看工作线程卡在哪里:如果卡在 socketRead,查下游接口和超时;卡在数据库连接,查连接池、慢 SQL 和事务;大量 BLOCKED 查锁竞争;大量 Future.get 查线程饥饿死锁;CPU 高则查计算和上下文切换。处理时要结合下游容量,不能盲目加线程,否则可能把数据库或第三方接口打得更慢。追问点:
- activeCount 满但 CPU 不高说明什么?
- 队列很大但没有拒绝是不是好事?
- 为什么线程池调大后吞吐不升反降?
- 如何判断是线程池瓶颈还是数据库瓶颈?
线程池源码执行链路怎么讲
标准回答:
ThreadPoolExecutor 提交任务后,先判断当前 workerCount 是否小于 corePoolSize,小于就创建 Worker 执行任务;否则尝试把任务放入 workQueue;入队成功后还要二次检查线程池状态,防止关闭后残留新任务;如果队列满了,再判断是否能创建到 maximumPoolSize;如果不能就执行拒绝策略。Worker 启动后会先执行 firstTask,然后在 runWorker 循环中不断通过 getTask 从队列取任务。getTask 会根据线程池状态、keepAliveTime、allowCoreThreadTimeOut 判断线程是等待任务还是退出回收。追问点:
- 为什么无界队列会让最大线程数基本失效?
getTask()为什么有时take,有时poll?- 为什么入队后还要二次检查线程池状态?
submit为什么可能隐藏异常?
关联知识点
| 知识点 | 说明 |
|---|---|
| 线程池总览 | 线程池学习路线 |
| 核心参数 | 七大参数和配置后果 |
| 执行流程原理 | execute 主流程 |
| 队列与拒绝策略 | 过载时如何排队、反压、拒绝 |
| 参数估算与监控 | 怎么按业务容量配置线程池 |
| 关闭与常见坑 | shutdown、submit 异常、饥饿死锁 |
| ThreadLocal | 线程复用下为什么必须 remove |
| JVM 排查工具 | jstack、jmap、jstat、dump 分析 |
本章小结
线程池是一个容量控制系统,不是一个简单异步工具。真正掌握线程池,要能把源码流程和生产现象对应起来:execute 决定任务怎么进入系统,Worker 决定线程怎么复用,getTask 决定线程怎么等待和退出,队列和拒绝策略决定过载时系统怎么保护自己,jstack 和监控指标决定你能不能在事故中找到真正瓶颈。
