线程池:execute() 为什么入队后要做二次检查?
一句话结论(30s)
execute() 入队后做二次检查是为了堵住两个竞态漏洞:任务入队瞬间可能恰好无线程消费、或线程池刚好 shutdown,因为入队是”成功”动作但后续执行并不被保证。关键设计是入队成功后重新读 ctl 状态,用 addWorker(null, false) 补一个消费线程、用 remove(command) + reject(command) 保证拒绝语义一致。权衡是用两次额外检查的极小开销,换任务”入队即有人消费、关闭即不被悬空”的语义确定性。
核心原理(2min)
先抛一个疑问:任务都已经 offer 进队列了,为什么还要回头再检查一遍?入队难道不是「成功」动作吗——带着它往下看。execute() 遵循「核心线程 → 队列 → 非核心线程 → 拒绝」四级提交策略:核心线程未满先 addWorker 建 core 线程,满了就 offer 入队,队满才建非核心线程,再满走拒绝策略。入队成功后会进入一个竞态窗口——最后一个 worker 可能刚因空闲超时退出(workerCount == 0),或 shutdown() 恰好发生(!isRunning)。所以入队后要重读 ctl 做两次检查:若已 shutdown 则从队列 remove 该任务并触发拒绝策略;若无线程则 addWorker(null, false) 补一个无 firstTask 的 worker 去消费队列,保证任务语义一致。
底层深入(5-10min)
execute() 的完整流程
下面是 JDK 中 execute() 的真实源码(含 Doug Lea 写的三段注释):
public void execute(Runnable command) {
Objects.requireNonNull(command, "command");
/*
* Proceed in 3 steps:
*
* 1. If fewer than corePoolSize threads are running, try to
* start a new thread with the given command as its first
* task. The call to addWorker atomically checks runState and
* workerCount, and so prevents false alarms that would add
* threads when it shouldn't, by returning false.
*
* 2. If a task can be successfully queued, then we still need
* to double-check whether we should have added a thread
* (because existing ones died since last checking) or that
* the pool shut down since entry into this method. So we
* recheck state and if necessary roll back the enqueuing if
* stopped, or start a new thread if there are none.
*
* 3. If we cannot queue task, then we try to add a new
* thread. If it fails, we know we are shut down or saturated
* and so reject the task.
*/
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);
}
真实源码里第二步的二次检查顺序是:先判 !isRunning(recheck)(若已 shutdown 则回滚入队并 reject),再判 workerCountOf(recheck) == 0(若无线程则补 addWorker(null, false))。两者用 else if 串联而非并列——两个竞态窗口在语义上互斥,同一时刻最多命中其一,避免重复处理。第一步的 addWorker(command, true) 内部会原子地同时校验 runState 和 workerCount,失败后再重读一次 c = ctl.get(),这正是 corePoolSize 二次检查的体现。
💭 思考:任务都已经
offer进队列了,为什么还要回头再检查一遍?入队成功难道不代表「任务有着落」吗?——一步步想:入队成功只意味着「任务进了队列这个盒子」,但「盒子有没有人在取」(workerCount是否为 0)和「盒子还在不在用」(池子是否已 shutdown)是另外两件事。入队这个动作和之后的检查之间隔着一个竞态窗口——恰好在这一瞬间,最后一个 worker 可能因空闲超时退出,shutdown()也可能刚好发生。如果不回头确认,任务就成了「既没人执行、也不会被拒绝」的孤儿。所以二次检查不是多此一举,而是把「入队成功」补齐成「入队必被消费、关闭必被拒绝」。
ctl 的位运算:一个 int 装下状态和线程数
二次检查之所以”读一次 ctl”就能同时拿到线程数和运行状态,——想没想过为什么「读一次 ctl」就能同时拿到两个值?靠的是 ctl 的位打包:
private final AtomicInteger ctl = new AtomicInteger(ctlOf(RUNNING, 0));
private static final int COUNT_BITS = Integer.SIZE - 3;
private static final int COUNT_MASK = (1 << COUNT_BITS) - 1;
// runState is stored in the high-order bits
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;
// Packing and unpacking ctl
private static int runStateOf(int c) { return c & ~COUNT_MASK; }
private static int workerCountOf(int c) { return c & COUNT_MASK; }
private static int ctlOf(int rs, int wc) { return rs | wc; }
private static boolean isRunning(int c) {
return c < SHUTDOWN;
}
COUNT_BITS = Integer.SIZE - 3 = 29,即低 29 位存 workerCount(上限约 5 亿),高 3 位存 runState。RUNNING = -1 << 29 在二进制里是 1110...0,是唯一负数状态,所以 isRunning(c) { return c < SHUTDOWN; } 用一次有符号整数比较就能判断——因为 runState 单调递增(RUNNING < SHUTDOWN < STOP < TIDYING < TERMINATED)。workerCountOf 用 & COUNT_MASK 取低 29 位、runStateOf 用 & ~COUNT_MASK 取高 3 位、ctlOf 用 rs | wc 拼回,这就是”一次 CAS 原子更新两个字段”的根因。
💭 思考:为什么要把「运行状态」和「线程数」塞进同一个 int,而不是用两个独立变量?——因为线程池的很多判断需要两者「同时一致」,比如「RUNNING 且 workerCount < corePoolSize」。如果分两个变量,一次判断要读两次、中间可能被别人改,竞态窗口就出现了;合成一个 int 后,一次 CAS 就能原子地读或改这两者,判断才不会「读到一半状态变了」。代价是 workerCount 被压缩到低 29 位、读写都要解包——这是拿「代码可读性」换「原子性」的经典取舍。
为什么需要入队后的二次检查?
检查一:!isRunning(recheck) —— shutdown 竞态
竞态场景:任务入队成功的一瞬间,线程池恰好被 shutdown()。任务悬空在队列里——既不会被执行(线程正在停止),也不会被拒绝(已经入队了)。试想:入队后若不回头检查,提交方已经拿到「入队成功」的结果,谁还会去处理这个悬空任务?它就成了一个既不被执行、也不被拒绝的「孤儿」任务。
源码里 if (! isRunning(recheck) && remove(command)) reject(command); 检测到这个状态后,先从队列 remove 掉该任务、再触发拒绝策略,保证”关闭即拒绝”的语义一致。注意这里判断的是 !isRunning 而非 == SHUTDOWN:STOP / TIDYING / TERMINATED 也都不应再接纳任务。
检查二:workerCount == 0 —— 无线程消费竞态
竞态场景:任务入队的瞬间,最后一个 worker 线程恰好因空闲超时退出(keepAliveTime 到期、allowCoreThreadTimeOut 配置)。此时队列不为空但无线程消费——任务永远不被执行。再问一层:队列里有任务,就一定有线程去消费吗?不一定,这里的竞态恰恰说明「入队成功」和「有人消费」是两回事。
addWorker(null, false) 创建一个新 worker(没有 firstTask),它从队列中拉取任务消费。这里传 null 且 core=false 是关键:任务已经在队列里,不需要 firstTask;且此刻可能没有核心线程存活,所以不能以 corePoolSize 为上限。
addWorker() 里的二次检查
可能你会追问:execute() 里那句「addWorker 原子地同时校验 runState 和 workerCount」是怎么做到的?二次检查不只在 execute() 里,addWorker() 内部同样做了双重校验,答案就在下面这段源码里:
private boolean addWorker(Runnable firstTask, boolean core) {
retry:
for (int c = ctl.get();;) {
// Check if queue empty only if necessary.
if (runStateAtLeast(c, SHUTDOWN)
&& (runStateAtLeast(c, STOP)
|| firstTask != null
|| workQueue.isEmpty()))
return false;
for (;;) {
if (workerCountOf(c)
>= ((core ? corePoolSize : maximumPoolSize) & COUNT_MASK))
return false;
if (compareAndIncrementWorkerCount(c))
break retry;
c = ctl.get(); // Re-read ctl
if (runStateAtLeast(c, SHUTDOWN))
continue retry;
// else CAS failed due to workerCount change; retry inner loop
}
}
boolean workerStarted = false;
boolean workerAdded = false;
Worker w = null;
try {
w = new Worker(firstTask);
final Thread t = w.thread;
if (t != null) {
final ReentrantLock mainLock = this.mainLock;
mainLock.lock();
try {
// Recheck while holding lock.
// Back out on ThreadFactory failure or if
// shut down before lock acquired.
int c = ctl.get();
if (isRunning(c) ||
(runStateLessThan(c, STOP) && firstTask == null)) {
if (t.getState() != Thread.State.NEW)
throw new IllegalThreadStateException();
workers.add(w);
workerAdded = true;
int s = workers.size();
if (s > largestPoolSize)
largestPoolSize = s;
}
} finally {
mainLock.unlock();
}
if (workerAdded) {
container.start(t);
workerStarted = true;
}
}
} finally {
if (! workerStarted)
addWorkerFailed(w);
}
return workerStarted;
}
第一层检查在自旋里:workerCountOf(c) >= 上限 直接返回 false,否则用 compareAndIncrementWorkerCount(c) CAS 自增线程数,失败就 c = ctl.get() 重读再试。第二层检查在持有 mainLock 后:isRunning(c) || (runStateLessThan(c, STOP) && firstTask == null) 再次确认此刻还能加线程——因为从 CAS 成功到拿锁之间,池子可能已经 shutdown,不重查就会把 worker 加进一个已关闭的池。这段代码同时解释了 execute() 里 addWorker 为何能”原子地同时校验 runState 和 workerCount”。
corePoolSize、workQueue、maxPoolSize 三者的配合
线程数增长策略:
新任务 → workerCount < corePoolSize? → 创建核心线程
→ 否 → 队列未满? → 入队
→ 否 → workerCount < maxPoolSize? → 创建非核心线程
→ 否 → 拒绝策略
经典的”非核心线程只在队列满时才创建”策略。
四种拒绝策略
| 策略 | 行为 |
|---|---|
| AbortPolicy(默认) | 抛 RejectedExecutionException |
| CallerRunsPolicy | 调用者线程自己执行任务(天然背压) |
| DiscardPolicy | 静默丢弃 |
| DiscardOldestPolicy | 丢弃队列中最旧的任务,重试提交 |
四种策略在 JDK 中的真实实现(ThreadPoolExecutor 内部静态类):
public static class CallerRunsPolicy implements RejectedExecutionHandler {
public CallerRunsPolicy() { }
public void rejectedExecution(Runnable r, ThreadPoolExecutor e) {
if (!e.isShutdown()) {
r.run();
}
}
}
public static class AbortPolicy implements RejectedExecutionHandler {
public AbortPolicy() { }
public void rejectedExecution(Runnable r, ThreadPoolExecutor e) {
throw new RejectedExecutionException("Task " + r.toString() +
" rejected from " +
e.toString());
}
}
public static class DiscardPolicy implements RejectedExecutionHandler {
public DiscardPolicy() { }
public void rejectedExecution(Runnable r, ThreadPoolExecutor e) {
}
}
public static class DiscardOldestPolicy implements RejectedExecutionHandler {
public DiscardOldestPolicy() { }
public void rejectedExecution(Runnable r, ThreadPoolExecutor e) {
if (!e.isShutdown()) {
e.getQueue().poll();
e.execute(r);
}
}
}
CallerRunsPolicy 有独特的”倒流”效果:任务提交速度太快 → 提交线程自己执行任务 → 提交速度自然降低 → 天然背压效果。DiscardOldestPolicy 里 e.getQueue().poll() 丢掉队头最旧任务后 e.execute(r) 重试,因为队列腾出了位置,重试大概率成功。注意除了 AbortPolicy 会直接抛异常外,其余三个都先判 !e.isShutdown()——已关闭的池子要静默丢弃而不是再执行或重试。
💭 思考:为什么 CallerRunsPolicy 能「天然背压」,而 AbortPolicy 不能?——背压的本质是让「提交速度」慢下来。AbortPolicy 直接抛异常,如果调用方不处理,提交风暴照样继续;CallerRunsPolicy 让调用线程亲自去跑这个任务,调用线程忙着执行、自然就没空再提交,提交速率被「执行时间」拖住,形成自反馈。Discard / DiscardOldest 则是用「丢任务」来保吞吐,各有各的代价。所以选择拒绝策略,本质是在「宁可失败」和「宁可减速」之间选立场。
线程数设计公式
CPU 密集型: 线程数 ≈ CPU 核数 + 1
IO 密集型: 线程数 ≈ CPU 核数 × (1 + IO等待时间/CPU计算时间)
《Java 并发编程实战》的原始公式。但注意这是对平台线程而言。虚拟线程(Java 21+)下,IO 密集型直接用 newVirtualThreadPerTaskExecutor(),无需计算线程数。
章末提问
- execute() 为什么入队后要做二次检查,不检查会怎样? —— 因为入队成功不代表一定有人消费、也不代表池子没关:入队瞬间可能最后一个 worker 刚超时退出(无线程消费),或恰好 shutdown(任务悬空),所以必须重读 ctl,要么
remove + reject、要么addWorker(null, false)补齐消费线程。 - 二次检查里的
addWorker(null, false)为什么不传任务、core 又为什么传 false? —— 因为任务已经躺在队列里,无需 firstTask;且此刻可能核心线程已全部退出,若按 core=true 受 corePoolSize 上限约束会拒绝创建,传 false 才能突破核心上限补一个消费者。 !isRunning(recheck)为什么不能写成== SHUTDOWN? —— 因为 STOP / TIDYING / TERMINATED 状态下池子同样不再接纳任务,而 isRunning 借助 runState 单调递增、RUNNING 是唯一负数这一事实,用一次比较就覆盖所有「不该再入队」的状态。- 检查一和检查二为什么用
else if串联,而不是两个独立 if? —— 因为两个竞态窗口语义互斥:shutdown 时不会再新建线程、无线程消费时池子仍在运行,同一时刻最多命中其一,串联能避免重复处理、逻辑更严谨。 - ctl 用一个 int 打包状态和线程数,好处和代价是什么? —— 好处是能一次 CAS 原子地同时更新两个字段、避免额外锁;代价是 workerCount 被压缩到低 29 位(上限约 5 亿),读写都要解包,代码可读性下降。