ThreadPoolExecutor 线程池源码精读
概述
ThreadPoolExecutor 是 JDK 线程池的实现核心:用一个 AtomicInteger ctl 同时打包运行状态(高 3 位)与工作线程数(低 29 位),通过 execute 四步策略、addWorker CAS 创建、Worker 内嵌 AQS、getTask 超时阻塞实现线程复用与池化。本文基于 OpenJDK 21 源码拆解完整执行链路。
一、ctl 的状态打包
1.1 位域设计
// java.util.concurrent.ThreadPoolExecutor
private final AtomicInteger ctl = new AtomicInteger(ctlOf(RUNNING, 0));
private static final int COUNT_BITS = Integer.SIZE - 3; // 29 位
private static final int COUNT_MASK = (1 << COUNT_BITS) - 1; // 低 29 位全 1
// 运行状态:高 3 位,低 29 位为 0
private static final int RUNNING = -1 << COUNT_BITS;
private static final int SHUTDOWN = 0 << COUNT_BITS;
private static final int STOP = 1 << COUNT_BITS;
private static final int TIDYING = 2 << COUNT_BITS;
private static final int TERMINATED = 3 << COUNT_BITS;
// 打包与解包
private static int runStateOf(int c) { return c & ~COUNT_MASK; } // 取高 3 位
private static int workerCountOf(int c) { return c & COUNT_MASK; } // 取低 29 位
private static int ctlOf(int rs, int wc) { return rs | wc; } // 合并用一个 int 同时保存两个变量,读写一次 CAS 即可完成"状态 + 数量"的原子切换,避免用两个字段时的复合竞态。
1.2 5 种运行状态
| 状态 | 数值(高 3 位) | 含义 | 是否接收新任务 | 是否处理队列任务 |
|---|---|---|---|---|
RUNNING | 111 | 正常运行 | 是 | 是 |
SHUTDOWN | 000 | 关闭(优雅) | 否 | 是(处理完存量) |
STOP | 001 | 立即停止 | 否 | 否(并中断所有线程) |
TIDYING | 010 | 收尾中(任务清空) | - | - |
TERMINATED | 011 | 已终止 | - | - |
状态单调递增:RUNNING → SHUTDOWN/STOP → TIDYING → TERMINATED,配合 CAS 保证状态转换安全。
private void advanceRunState(int targetState) {
for (;;) {
int c = ctl.get();
if (runStateAtLeast(c, targetState) ||
ctl.compareAndSet(c, ctlOf(targetState, workerCountOf(c))))
break;
}
}二、execute(Runnable) 四步提交策略
public void execute(Runnable command) {
if (command == null) throw new NullPointerException();
int c = ctl.get();
// ① 核心线程数未满 → 创建核心线程执行
if (workerCountOf(c) < corePoolSize) {
if (addWorker(command, true))
return;
c = ctl.get(); // 创建失败(如已达上限)→ 重新读取
}
// ② 运行中 → 尝试入队
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); // 没有线程 → 补一个(即使队列有任务也要有消费线程)
}
// ③ 非运行状态或入队失败 → 尝试创建非核心线程
else if (!addWorker(command, false))
reject(command); // ④ 都失败 → 拒绝策略
}提交策略总结:
1. 核心线程有空位 → 直接建线程执行任务
2. 否则入队(入队后双重检查状态与线程数)
3. 否则尝试建非核心线程(超出队列容量)
4. 全部失败 → 执行拒绝策略核心线程的创建是惰性的:任务到来才建,不是启动就建(
prestartCoreThread可预热,见第七节)。
三、addWorker(Runnable, boolean) 创建线程
private boolean addWorker(Runnable firstTask, boolean core) {
retry:
for (;;) {
int c = ctl.get();
int rs = runStateOf(c);
// 状态检查:非 RUNNING 且(非 SHUTDOWN 或 firstTask 非空或队列非空)→ 拒绝
if (rs >= SHUTDOWN &&
!(rs == SHUTDOWN && firstTask == null && !workQueue.isEmpty()))
return false;
for (;;) {
int wc = workerCountOf(c);
if (wc >= CAPACITY ||
wc >= (core ? corePoolSize : maximumPoolSize))
return false; // 数量超限
if (compareAndIncrementWorkerCount(c)) // CAS 增加 worker 计数
break retry; // 计数成功 → 跳出
c = ctl.get(); // CAS 失败 → 重读
if (runStateOf(c) != rs) continue retry; // 状态变了 → 重新外层循环
}
}
boolean workerStarted = false;
boolean workerAdded = false;
Worker w = null;
try {
w = new Worker(firstTask); // ① 新建 Worker(自带 AQS)
final Thread t = w.thread;
if (t != null) {
final ReentrantLock mainLock = this.mainLock;
mainLock.lock(); // ② workers 集合操作加锁
try {
int rs = runStateOf(ctl.get());
if (rs < SHUTDOWN ||
(rs == SHUTDOWN && firstTask == null)) {
if (t.isAlive()) throw new IllegalThreadStateException();
workers.add(w); // ③ 加入 workers 集合
workerAdded = true;
}
} finally {
mainLock.unlock();
}
if (workerAdded) {
t.start(); // ④ 启动线程
workerStarted = true;
}
}
} finally {
if (!workerStarted)
addWorkerFailed(w); // 启动失败 → 回滚计数并清理
}
return workerStarted;
}addWorker 外层自旋 CAS 增加 workerCount,成功后加锁操作 workers(HashSet<Worker>,需 mainLock 保护),最后 t.start()。
四、Worker:线程池的最小工作单元
// Worker 自身继承 AQS!使用互斥锁语义(不可重入)
private final class Worker extends AbstractQueuedSynchronizer implements Runnable {
final Thread thread; // 真正执行任务的线程(ThreadFactory 创建)
Runnable firstTask; // 首个任务(可能是 null,表示只从队列取)
volatile long completedTasks; // 完成的任务数(统计用)
Worker(Runnable firstTask) {
setState(-1); // 初始 AQS state = -1,禁止中断
this.firstTask = firstTask;
this.thread = getThreadFactory().newThread(this);
}
public void run() { runWorker(this); }
// —— AQS 重写:互斥(不可重入)锁 ——
protected boolean tryAcquire(int unused) {
if (compareAndSetState(0, 1)) { setExclusiveOwnerThread(Thread.currentThread()); return true; }
return false;
}
protected boolean tryRelease(int unused) { setExclusiveOwnerThread(null); setState(0); return true; }
public boolean tryLock() { return tryAcquire(1); }
public void unlock() { release(1); }
public boolean isLocked() { return isHeldExclusively(); }
}Worker 复用 AQS 作为不可重入的自身锁:shutdownNow() 中断线程前要 tryLock() 成功(即 worker 不在 runWorker 中被锁住),防止中断正在执行任务的线程导致任务执行一半被暴力打断。
五、runWorker 主循环
final void runWorker(Worker w) {
Thread wt = Thread.currentThread();
Runnable task = w.firstTask;
w.firstTask = null;
w.unlock(); // 允许中断(state 从 -1 → 0)
boolean completedAbruptly = true;
try {
// 核心循环:先执行 firstTask,再不断从队列取任务
while (task != null || (task = getTask()) != null) {
w.lock(); // 执行任务期间上锁 → shutdownNow 无法中断
try {
if ((runStateAtLeast(ctl.get(), STOP) ||
(Thread.interrupted() && runStateAtLeast(ctl.get(), STOP))) &&
!wt.isInterrupted())
wt.interrupt(); // STOP 状态下保证中断
try {
task.run(); // 真正执行任务(beforeExecute → task → afterExecute)
completedAbruptly = false;
} finally {
task = null;
w.completedTasks++;
w.unlock(); // 任务执行完解锁
}
} catch (Throwable x) {
// 任务抛异常 → 线程退出(completedAbruptly = true)
} finally {
afterExecute(task, thrown);
}
}
completedAbruptly = false;
} finally {
processWorkerExit(w, completedAbruptly); // 线程退出清理(可能补线程)
}
}任务执行期间 w.lock() 的意义:shutdownNow() 遍历 workers 逐个中断时,正在执行任务的 worker 拿不到锁,中断会等到任务执行完释放锁;但 STOP 状态会强制给线程打中断标记。
// processWorkerExit 的部分逻辑
private void processWorkerExit(Worker w, boolean completedAbruptly) {
if (completedAbruptly) // 异常退出 → 先回滚 worker 计数
decrementWorkerCount();
// ... 从 workers 移除、更新 completedTaskCount、tryTerminate()
// 补偿:若仍 RUNNING 且线程数不足 corePoolSize → 补线程
if (runStateLessThan(ctl.get(), STOP)) {
if (!completedAbruptly) {
int min = allowCoreThreadTimeOut ? 0 : corePoolSize;
if (min == 0 && !workQueue.isEmpty()) min = 1;
if (workerCountOf(ctl.get()) >= min) return;
}
addWorker(null, false); // 补一个线程继续服务队列
}
}六、getTask() 的阻塞与超时
private Runnable getTask() {
boolean timedOut = false;
for (;;) {
int c = ctl.get();
int rs = runStateOf(c);
// ① SHUTDOWN 且队列空 / STOP → 不取任务,返回 null(线程退出)
if (rs >= SHUTDOWN && (rs >= STOP || workQueue.isEmpty())) {
decrementWorkerCount();
return null;
}
int wc = workerCountOf(c);
boolean timed = allowCoreThreadTimeOut || wc > corePoolSize;
// ② 超时策略:线程数超核心或允许核心超时 → 用 poll(timeout)
if ((wc > maximumPoolSize || (timed && timedOut))
&& (wc > 1 || workQueue.isEmpty())) {
if (compareAndDecrementWorkerCount(c)) return null;
continue;
}
try {
Runnable r = timed ?
workQueue.poll(keepAliveTime, TimeUnit.NANOSECONDS) : // 超时阻塞
workQueue.take(); // 永久阻塞
if (r != null) return r;
timedOut = true; // poll 超时返回 null → 下一轮判定退出
} catch (InterruptedException retry) {
timedOut = false; // 被中断 → 重新循环(不退出)
}
}
}poll(timeout) 与 take() 的区别:
workQueue.poll(timeout):
队列为空 → 阻塞最多 timeout 后返回 null
用于:核心线程超时淘汰(allowCoreThreadTimeOut=true)或
非核心线程空闲超时(wc > corePoolSize)回收
workQueue.take():
队列为空 → 永久阻塞(等生产者放入)
用于:核心线程(默认不回收)无限等待任务七、四种拒绝策略与线程预热
7.1 拒绝策略
reject(command) 委托给 RejectedExecutionHandler handler,四种内置实现:
// ① 默认:AbortPolicy —— 直接抛异常
public void rejectedExecution(Runnable r, ThreadPoolExecutor e) {
throw new RejectedExecutionException("Task " + r.toString() + " rejected from " + e.toString());
}
// ② CallerRunsPolicy —— 由调用线程自己执行(背压)
public void rejectedExecution(Runnable r, ThreadPoolExecutor e) {
if (!e.isShutdown()) r.run(); // 不进入线程池,直接调用线程跑
}
// ③ DiscardPolicy —— 静默丢弃
public void rejectedExecution(Runnable r, ThreadPoolExecutor e) { }
// ④ DiscardOldestPolicy —— 丢弃队列头(最旧)任务,再提交
public void rejectedExecution(Runnable r, ThreadPoolExecutor e) {
if (!e.isShutdown()) {
e.getQueue().poll(); // 丢最旧
e.execute(r); // 重试新任务
}
}拒绝策略适用场景:
AbortPolicy:默认,任务不可丢(会抛异常暴露问题)
CallerRunsPolicy:需要节流降速(调用线程执行,天然背压)
DiscardPolicy:允许丢任务(日志上报等可丢场景)
DiscardOldestPolicy:任务有新鲜度要求(丢弃过期任务)7.2 线程预热
public boolean prestartCoreThread() {
return workerCountOf(ctl.get()) < corePoolSize &&
addWorker(null, true); // 提前创建核心线程(无首任务,只取队列)
}
public int prestartAllCoreThreads() {
int n = 0;
while (addWorker(null, true)) // 循环创建直到核心数满
++n;
return n;
}核心线程默认惰性创建,高并发尖峰场景可 prestartAllCoreThreads() 预热,避免第一波任务全部走"建线程"慢路径。
八、实现要点
ThreadPoolExecutor 核心:
ctl 位打包:高 3 位状态 + 低 29 位线程数,单次 CAS 原子切换
状态机:RUNNING → SHUTDOWN/STOP → TIDYING → TERMINATED
execute 四步:建核心 → 入队 → 建非核心 → 拒绝
Worker:线程载体 + 自身 AQS 互斥锁(防中断正在执行的任务)
getTask:poll(timeout) 超时回收 vs take() 永久阻塞
拒绝策略:抛异常 / 调用者执行 / 丢弃 / 丢最旧
预热:prestartCoreThread / prestartAllCoreThreads
常见陷阱:
核心线程不回收(需 allowCoreThreadTimeOut)
队列用无界队列 → 非核心线程永不创建,拒绝策略失效
shutdownNow 中断正在执行的任务 → Worker AQS 锁保护
任务抛异常线程退出 → processWorkerExit 自动补线程