同事问了个我答不上来的问题
2 月中旬,组里做短信发送模块的改造,同事在配置线程池时问我:corePoolSize 用完之后,是立刻扩容到 maximumPoolSize,还是先往队列里塞?
我当时脱口而出"先扩容到 max"。说完自己就心虚了,因为印象里看过"队列满了才会扩"的说法。答不上来的问题最丢人,当晚我把 JDK 8 的 ThreadPoolExecutor 源码读了一遍,顺便做了个实验验证。
先用一个 demo 把行为逼出来
参数故意设得很小:核心 2、最大 4、队列容量 2,任务是睡 10 秒的空任务,方便观察中间状态。
public class PoolStepDemo {
public static void main(String[] args) {
ThreadPoolExecutor pool = new ThreadPoolExecutor(
2, // corePoolSize
4, // maximumPoolSize
30L, TimeUnit.SECONDS,
new ArrayBlockingQueue<>(2), // 有界队列,容量 2
new NamedThreadFactory("sms-pool"),
new ThreadPoolExecutor.AbortPolicy());
for (int i = 1; i <= 7; i++) {
final int no = i;
try {
pool.execute(() -> {
try { Thread.sleep(10000); } catch (InterruptedException e) { }
System.out.println(Thread.currentThread().getName() + " 执行完 task-" + no);
});
System.out.printf("提交 task-%d 成功 | poolSize=%d, queue=%d%n",
no, pool.getPoolSize(), pool.getQueue().size());
} catch (RejectedExecutionException e) {
System.out.printf("提交 task-%d 被拒绝 | poolSize=%d, queue=%d%n",
no, pool.getPoolSize(), pool.getQueue().size());
}
}
pool.shutdown();
}
}
输出:
提交 task-1 成功 | poolSize=1, queue=0
提交 task-2 成功 | poolSize=2, queue=0
提交 task-3 成功 | poolSize=2, queue=1
提交 task-4 成功 | poolSize=2, queue=2
提交 task-5 成功 | poolSize=3, queue=2
提交 task-6 成功 | poolSize=4, queue=2
提交 task-7 被拒绝 | poolSize=4, queue=2
结论很清楚:核心线程满了之后,任务先入队列;队列满了,才会创建超过核心数的线程;线程数到 maximumPoolSize 后还塞不下,才走拒绝策略。我当时的答案是错的。
execute 的三段判断
源码就二十多行,三段 if 分别对应上面的三个阶段:
public void execute(Runnable command) {
if (command == null)
throw new NullPointerException();
int c = ctl.get();
// 第一段:当前工作线程数 < corePoolSize,直接新建核心线程
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);
}
这里有个容易看漏的细节:ctl 是一个 AtomicInteger,高 3 位存线程池状态(RUNNING / SHUTDOWN / STOP / TIDYING / TERMINATED),低 29 位存工作线程数。两个状态放在一个 int 里,是为了用一次 CAS 同时更新,避免加锁。
private static final int COUNT_BITS = Integer.SIZE - 3; // 29
private static final int CAPACITY = (1 << COUNT_BITS) - 1; // 约 5.3 亿
private static int runStateOf(int c) { return c & ~CAPACITY; }
private static int workerCountOf(int c) { return c & CAPACITY; }
第二段里的双重检查值得留意:入队成功之后线程池可能刚好被 shutdown() 了,所以要重新读一次状态,发现不在运行就把任务从队列里 remove 掉再拒绝。这就是为什么 shutdown() 之后队列里的任务还会继续执行完,但新提交的任务会被拒绝。
addWorker:线程是这么被创建出来的
两段 addWorker 调用只有一个布尔参数不同:true 表示按 corePoolSize 校验上限,false 表示按 maximumPoolSize 校验。这就是"核心"和"非核心"唯一的差别,线程本身没有任何标记区分,"核心线程"只是一个数量概念。
private boolean addWorker(Runnable firstTask, boolean core) {
retry:
for (;;) {
int c = ctl.get();
int rs = runStateOf(c);
// 状态检查:SHUTDOWN 之后不再接新任务,但允许处理队列里剩下的
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; // 到上限了,返回 false
if (compareAndIncrementWorkerCount(c)) // CAS 增加 workerCount
break retry;
c = ctl.get();
if (runStateOf(c) != rs)
continue retry;
// CAS 失败说明有别人抢先改了计数,重来
}
}
boolean workerStarted = false;
boolean workerAdded = false;
Worker w = null;
try {
w = new Worker(firstTask); // 把任务包装成 Worker
final Thread t = w.thread;
if (t != null) {
final ReentrantLock mainLock = this.mainLock;
mainLock.lock();
try {
int rs = runStateOf(ctl.get());
if (rs < SHUTDOWN ||
(rs == SHUTDOWN && firstTask == null)) {
if (t.isAlive())
throw new IllegalThreadStateException();
workers.add(w); // HashSet<Worker>,要加锁
int s = workers.size();
if (s > largestPoolSize)
largestPoolSize = s;
workerAdded = true;
}
} finally {
mainLock.unlock();
}
if (workerAdded) {
t.start(); // 真正启动线程
workerStarted = true;
}
}
} finally {
if (! workerAdded)
addWorkerFailed(w); // 回滚计数并移除
}
return workerStarted;
}
注意 workers 是个 HashSet,非线程安全,所以操作它要先拿 mainLock。而 workerCount 的增减用 CAS。这个分工是 AQS 之外的另一套思路:能用原子操作解决的就不加锁。
还有一处我以前没注意:workers.add 和 t.start() 都放在 mainLock 之外(start 在 unlock 之后)。Doug Lea 注释里说明这是为了减小锁的持有范围。
getTask:线程为什么不会死,又为什么会死
线程启动后跑的是 runWorker,一个循环:先执行自己带的 firstTask,然后不断从队列里取任务。
final void runWorker(Worker w) {
Runnable task = w.firstTask;
w.firstTask = null;
w.unlock(); // Worker 继承 AQS,unlock 允许被中断
boolean completedAbruptly = true;
try {
while (task != null || (task = getTask()) != null) {
w.lock();
// 如果线程池进入 STOP,确保当前线程被中断
if ((runStateAtLeast(ctl.get(), STOP) ||
(Thread.interrupted() && runStateAtLeast(ctl.get(), STOP))) &&
!w.thread.isInterrupted())
w.thread.interrupt();
try {
beforeExecute(w.thread, task); // 钩子方法,留给子类
Throwable thrown = null;
try {
task.run();
} catch (RuntimeException x) {
thrown = x; throw x; // 异常继续往外抛
} finally {
afterExecute(task, thrown);
}
} finally {
task = null;
w.completedTasks++;
w.unlock();
}
}
completedAbruptly = false;
} finally {
processWorkerExit(w, completedAbruptly);
}
}
线程能不能活下去,全看 getTask():
private Runnable getTask() {
boolean timedOut = false;
for (;;) {
int c = ctl.get();
int rs = runStateOf(c);
// 线程池 STOP 了,或者 SHUTDOWN 且队列已空 → 回收线程
if (rs >= SHUTDOWN && (rs >= STOP || workQueue.isEmpty())) {
decrementWorkerCount();
return null; // 返回 null 意味着这个线程退出循环
}
int wc = workerCountOf(c);
boolean timed = allowCoreThreadTimeOut || wc > corePoolSize;
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 超时了,下一轮判断要不要回收
} catch (InterruptedException retry) {
timedOut = false;
}
}
}
这段代码解释了两件我一直模糊的事:
keepAliveTime默认只对超出corePoolSize的那部分线程生效。timed这个变量决定了用poll(带超时)还是take(永久阻塞)。核心线程走take,所以永远不会被回收,除非你显式开了allowCoreThreadTimeOut(true)。- 线程数的收缩是被动的:没有定时任务去扫描,而是线程自己在
poll超时之后,下一轮循环里发现wc > corePoolSize且自己超时了,才 CAS 减计数并返回 null 退出。所以线程池从 4 缩回 2,最快也要等一个keepAliveTime。
processWorkerExit 里还有个细节:如果任务是抛异常退出的(completedAbruptly = true),它会直接调 addWorker(null, false) 补一个线程上来。也就是说线程池里的线程不会因为任务抛异常就少一个。
拒绝策略的触发条件
从 execute 的第三段可以看出,reject(command) 只在"队列 offer 失败 且 addWorker 失败"时执行。JDK 8 内置四种:
| 策略 | 行为 | 我们用在 |
|---|---|---|
| AbortPolicy(默认) | 抛 RejectedExecutionException | 核心的支付、下单 |
| CallerRunsPolicy | 让提交任务的线程自己跑 | 短信、日志这种可降级的 |
| DiscardPolicy | 静默丢弃 | 没用过,太危险 |
| DiscardOldestPolicy | 丢掉队列头,再试一次提交 | 没用过 |
CallerRunsPolicy 其实是个挺巧妙的设计:提交方(通常是 Tomcat 工作线程)被迫自己执行任务,这段时间它没法提交新任务,相当于给上游一个自然的反压。我们短信模块就用的它,压测时 QPS 从 1200 降到 780,但没有一条短信丢失。
我们的最终配置
之前项目里到处是 Executors.newFixedThreadPool(10),被组长点了两次名。原因是它内部用的 LinkedBlockingQueue 没指定容量,默认是 Integer.MAX_VALUE,任务堆积时队列永远填不满,maximumPoolSize 形同虚设,最后 OOM。同理 newCachedThreadPool 的 maximumPoolSize 是 Integer.MAX_VALUE,来多少任务建多少线程。
@Bean("smsExecutor")
public ThreadPoolExecutor smsExecutor() {
int core = Runtime.getRuntime().availableProcessors(); // 4 核机器 → 4
return new ThreadPoolExecutor(
core,
core * 2,
60L, TimeUnit.SECONDS,
new ArrayBlockingQueue<>(500), // 必须指定容量
new NamedThreadFactory("sms"), // 自定义线程名,排查时太重要了
new ThreadPoolExecutor.CallerRunsPolicy());
}
自定义线程工厂是吃过亏才加上去的。之前线上线程池里的线程都叫 pool-7-thread-3,jstack 出来一堆这种名字,根本分不清是哪个业务的。改成 sms-1、order-2 之后,一眼就能看出谁在堆积。
另外加了监控,每分钟上报一次 getPoolSize()、getQueue().size()、getCompletedTaskCount(),队列长度超过 80% 就告警。
下篇预告
这篇先把《ThreadPoolExecutor 源码解析:execute 之后发生了什么》里的坑列了,下一篇写我们当时是怎么在线上工程里真正落地的——包括那次让领导拍桌的故障复盘。