为什么需要线程池
如果每个任务都 new Thread() 去执行:创建和销毁线程本身就有开销,而且销毁后线程无法复用,CPU 缓存也就失效了。线程池的核心思路是 复用线程:预先创建好一批线程,任务来了就交给空闲线程去跑,避免反复创建销毁线程。
判断一个优秀线程池的标准:
- 复用:任务交给已有的空闲线程,而不是每次都新建
- 控制并发数:避免无限制创建线程耗尽资源
- 管理生命周期:空闲线程可以被回收,池可以被优雅关闭
核心参数
ThreadPoolExecutor 有 7 个构造参数 1:
public ThreadPoolExecutor(
int corePoolSize, // 核心线程数
int maximumPoolSize, // 最大线程数
long keepAliveTime, // 空闲线程存活时间
TimeUnit unit, // keepAliveTime 的时间单位
BlockingQueue<Runnable> workQueue, // 任务队列
ThreadFactory threadFactory, // 线程工厂
RejectedExecutionHandler handler // 拒绝策略
) {
// ...
}ThreadPoolExecutor.java
| 参数 | 作用 |
|---|---|
corePoolSize | 核心线程数,即使空闲也不会被回收(除非 allowCoreThreadTimeOut) |
maximumPoolSize | 最大线程数,线程池能创建线程的上限 |
keepAliveTime | 超过核心线程数部分的线程,空闲超过该时间就会被回收 |
unit | keepAliveTime 的时间单位 |
workQueue | 任务队列,当线程数达到 corePoolSize 后,新任务先入队 |
threadFactory | 线程工厂,用于创建线程,可自定义线程名/优先级/是否守护线程 |
handler | 拒绝策略,线程池无法接收新任务时的处理方式 |
任务来了怎么走:如果线程数少于 corePoolSize,直接创建新线程执行;如果已经达到 corePoolSize,任务放入 workQueue 等待;如果队列满了,再创建线程直到 maximumPoolSize;如果线程数已经达到 maximumPoolSize 且队列也满了,就交给拒绝策略处理。
工作队列 workQueue
workQueue 决定了「任务排队」的规则,它直接影响了线程池如何增长线程。常用的队列有:
| 队列 | 类型 | 是否有界 | 特点 |
|---|---|---|---|
ArrayBlockingQueue | 基于数组 | 有界(初始化指定) | 适合对任务堆积数有明确上限、资源可控的场景 |
LinkedBlockingQueue | 基于链表 | 默认无界,可传容量 | 可以无界;newFixedThreadPool 默认用它 |
SynchronousQueue | 移交队列 | 容量为 0 | 本身不存任务,每次 put 必须等一个 take 配对,用于直接移交 |
PriorityBlockingQueue | 优先级队列 | 默认可扩容 | 任务按优先级出队 |
DelayQueue | 延迟队列 | 无界 | 任务到期才能出队,常用于调度 |
队列类型如何影响线程增长? 上一节的执行流程里有一个隐含前提:队列是否容易满,决定了线程数会不会往 maximumPoolSize 涨。
- 用无界队列(如默认的
LinkedBlockingQueue):任务永远进得了队,队列永远不会「满」。于是execute永远不会走到第 3 步的addWorker(command, false),线程数最多涨到corePoolSize就不涨了,maximumPoolSize形同虚设。缺点是任务多了会无限堆积,可能 OOM。 - 用
SynchronousQueue(容量为 0):任务根本没法排队,offer永远失败。于是一来任务就直接走第 3 步去建非核心线程,线程数能一路涨到maximumPoolSize。这正是newCachedThreadPool的行为——线程可随时创建,闲下来 60s 回收。
所以「有界队列 + 明确的拒绝策略」通常是生产环境的稳妥组合:它既能限制堆积量、防止 OOM,又能让 maximumPoolSize 真正生效——当队列满了、线程也满了,就用拒绝策略兜底。
ctl:打包线程池状态
线程池需要同时维护「线程池运行状态」和「工作线程数量」两个信息。为了把它们放在一个变量里原子更新,JDK 用一个 AtomicInteger 类型的 ctl 打包了这两个字段:
private final AtomicInteger ctl = new AtomicInteger(ctlOf(RUNNING, 0));
private static final int COUNT_BITS = Integer.SIZE - 3; // 32 - 3 = 29
// 取低 29 位作为工作线程数,高 3 位作为运行状态
private static final int CAPACITY = (1 << COUNT_BITS) - 1;
private static final int RUNNING = -1 << COUNT_BITS; // 高 3 位 111
private static final int SHUTDOWN = 0 << COUNT_BITS; // 高 3 位 000
private static final int STOP = 1 << COUNT_BITS; // 高 3 位 001
private static final int TIDYING = 2 << COUNT_BITS; // 高 3 位 010
private static final int TERMINATED = 3 << COUNT_BITS; // 高 3 位 011
private static int runStateOf(int c) { return c & ~CAPACITY; } // 取高 3 位
private static int workerCountOf(int c) { return c & CAPACITY; } // 取低 29 位ThreadPoolExecutor.java
- 高 3 位:线程池的运行状态
- 低 29 位:工作线程的数量(上限
CAPACITY = 2^29 - 1)
把状态和数量压缩进一个 int,就能用一次 CAS 同时更新「状态」和「线程数」,保证两者的一致性和原子性。
线程池的五种状态
状态之间的大小关系是 RUNNING < SHUTDOWN < STOP < TIDYING < TERMINATED,所以源码里用 runStateAtLeast(c, STOP) 这种数值比较来判断状态:
| 状态 | 含义 |
|---|---|
RUNNING | 正常运行:接收新任务,处理队列里的任务 |
SHUTDOWN | 调用 shutdown() 后:不接收新任务,但会处理队列里已有的任务 |
STOP | 调用 shutdownNow() 后:不接收新任务,不处理队列任务,并中断正在执行的任务 |
TIDYING | 所有任务已结束、工作线程数为 0,即将执行 terminated() 钩子方法 |
TERMINATED | terminated() 已执行完毕,线程池彻底终止 |
SHUTDOWN 和 STOP 的关键区别在于:SHUTDOWN 会继续把队列里的任务跑完,而 STOP 直接丢弃队列任务并中断正在执行的线程。
execute() 的流程
提交任务的核心方法是 execute(Runnable command):
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(); // addWorker 失败(比如状态已变),重新读取
}
// 第二步:线程数 >= corePoolSize,尝试把任务放入队列
if (isRunning(c) && workQueue.offer(command)) {
int recheck = ctl.get();
// 二次检查①:入队后线程池被关闭 → 移除并拒绝,不让任务烂在队列里
if (!isRunning(recheck) && remove(command))
reject(command);
// 二次检查②:工作线程为 0(核心线程刚被超时回收/退出)→ 补一个线程去取任务
// 注意 firstTask 传 null:任务已在队列里,传 command 会被执行两次
else if (workerCountOf(recheck) == 0)
addWorker(null, false);
}
// 第三步:队列已满,尝试新建非核心线程执行,失败则拒绝
else if (!addWorker(command, false))
reject(command);
}ThreadPoolExecutor.java
整个逻辑可以概括为三步:核心线程 → 入队 → 非核心线程 → 拒绝。
其中第二步的「二次检查」很关键:任务入队成功不代表万事大吉。入队之后到检查之前,线程池可能已经被其他线程关闭了,此时必须把任务从队列里移除再交给拒绝策略;同时还要防止因核心线程为空(比如被超时回收或正在退出)而出现「队列里有任务却没人跑」的情况——这时要补一个线程,它的 firstTask 传 null,因为它没有任何初始任务,直接进入 getTask() 循环去取队列里那个已经在等的任务;若把 command 传进去,这个任务就会被执行两次。
addWorker:创建线程
addWorker 负责真正创建一个线程,并通过一个布尔参数 core 区分创建的是核心线程还是非核心线程:
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);
// 线程数检查:超过上限,或超过 core/maximum 对应的数量
if (wc >= CAPACITY ||
wc >= (core ? corePoolSize : maximumPoolSize))
return false;
// CAS 增加工作线程数,成功后跳出外层循环
if (compareAndIncrementWorkerCount(c))
break retry;
c = ctl.get();
if (runStateOf(c) != rs)
continue retry; // 状态变了,重新走外层检查
// 否则是线程数被别的线程改了,继续内层循环重试
}
}
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 {
int rs = runStateOf(ctl.get());
// 再检查一次状态,防止锁竞争期间线程池被关闭
if (rs < SHUTDOWN ||
(rs == SHUTDOWN && firstTask == null)) {
if (t.isAlive())
throw new IllegalThreadStateException();
workers.add(w);
int s = workers.size();
if (s > largestPoolSize)
largestPoolSize = s;
workerAdded = true;
}
} finally {
mainLock.unlock();
}
if (workerAdded) {
t.start(); // 真正启动线程,去执行 runWorker
workerStarted = true;
}
}
} finally {
if (!workerStarted)
addWorkerFailed(w); // 启动失败则回滚:移除线程并递减计数
}
return workerStarted;
}ThreadPoolExecutor.java
addWorker 分为两段:
- 前段:用 CAS 更新
ctl,把工作线程数加 1。这里的双重循环是为了同时应对「状态变化」和「线程数被其他线程抢占」两种竞争。 - 后段:创建
Worker、加进workers集合,然后启动线程。启动失败时要调用addWorkerFailed回滚,把线程从集合移除并把计数减回去。
Worker:线程与任务的捆绑
每个工作线程都被包装成一个 Worker 对象,它同时是「一个线程」和「一个任务」的载体:
private final class Worker extends AbstractQueuedSynchronizer
implements Runnable {
final Thread thread; // 这个 Worker 真正对应的线程
Runnable firstTask; // 创建时携带的第一个任务,可能为 null
volatile long completedTasks; // 已完成的任务数
Worker(Runnable firstTask) {
setState(-1); // 禁止中断,直到 runWorker 执行
this.firstTask = firstTask;
this.thread = getThreadFactory().newThread(this);
}
public void run() {
runWorker(this); // 线程启动后进入这个循环
}
}ThreadPoolExecutor.java
为什么 Worker 要继承 AQS?因为在 runWorker 里会对线程加锁(w.lock() / w.unlock())。这个锁不是为了互斥访问临界区,而是为了 标记「当前 Worker 是否正在执行任务」:
- Worker 在执行任务前
lock(),执行完后unlock() - 关闭线程池(
shutdown())时只中断 空闲 的 Worker(没有加锁的 Worker),正在执行任务的 Worker 不会被中断
这样就能做到「只中断空闲线程,不打断正在跑任务的线程」。
runWorker:死循环取任务
线程启动后进入 runWorker,这是一个不断「取任务 → 执行」的死循环,也是线程复用的核心:
final void runWorker(Worker w) {
Thread wt = Thread.currentThread();
Runnable task = w.firstTask;
w.firstTask = null;
w.unlock(); // 解除构造时的 -1 状态,允许后续中断
boolean completedAbruptly = true;
try {
// 先执行 firstTask,之后不断从队列取任务。取不到(返回 null)就退出
while (task != null || (task = getTask()) != null) {
w.lock();
if ((runStateAtLeast(ctl.get(), STOP) ||
(Thread.interrupted() &&
runStateAtLeast(ctl.get(), STOP))) &&
!wt.isInterrupted())
wt.interrupt();
try {
beforeExecute(wt, task);
try {
task.run();
afterExecute(task, null);
} catch (Throwable ex) {
afterExecute(task, ex); // 钩子能回调到异常
throw ex;
}
} finally {
task = null;
w.completedTasks++;
w.unlock();
}
}
completedAbruptly = false;
} finally {
processWorkerExit(w, completedAbruptly);
}
}ThreadPoolExecutor.java
beforeExecute/afterExecute是两个 可重写钩子,允许在任务前后做额外处理(比如埋点、记录日志)- 循环退出后调用
processWorkerExit,清理 Worker 并可能补位(如果队列里还有任务)
getTask:取任务与超时回收
getTask 决定线程是「阻塞等待任务」还是「超时后退出」。空闲线程的回收就发生在这里:
private Runnable getTask() {
boolean timedOut = false; // 上次 poll 是否超时
for (;;) {
int c = ctl.get();
int rs = runStateOf(c);
// 停止状态,或队列为空 + 非 RUNNING → 线程退出
if (rs >= SHUTDOWN && (rs >= STOP || workQueue.isEmpty())) {
decrementWorkerCount();
return null;
}
int wc = workerCountOf(c);
// 是否参与超时回收:
// 1. 允许核心线程超时,或
// 2. 当前线程数 > corePoolSize(多余的线程才会被回收)
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; // 被中断则重置,重新循环
}
}
}ThreadPoolExecutor.java
keepAliveTime 的作用在这里体现:
- 线程数 不超过
corePoolSize时,timed = false,用workQueue.take()永久阻塞 等任务,这类线程不会被回收 - 线程数 超过
corePoolSize时,timed = true,用workQueue.poll(keepAliveTime)超时等待;一旦超时(timedOut = true),线程就退出循环、被回收
所以「超过核心线程数的那些线程在空闲 keepAliveTime 后会被回收」。如果设置了 allowCoreThreadTimeOut = true,核心线程也会被回收。
拒绝策略
当任务无法提交(线程数达到 maximumPoolSize 且队列已满,或线程池已关闭)时,会交给 RejectedExecutionHandler 处理。JDK 内置了四种策略:
| 策略 | 行为 | 风险 |
|---|---|---|
AbortPolicy | 默认策略,直接抛出 RejectedExecutionException | 调用方需自行捕获异常 |
CallerRunsPolicy | 让 提交任务的线程 自己执行该任务 | 提交线程被占用,可能拖慢流程 |
DiscardPolicy | 静默丢弃新任务 | 任务可能丢失,无任何提示 |
DiscardOldestPolicy | 丢弃队列中 最老 的任务,然后重试提交新任务 | 可能丢任务,适合时效性高的场景 |
CallerRunsPolicy 是一种「限流」手段:当线程池跑不动时,让提交者自己跑,一方面起到背压作用,另一方面任务不会被丢弃。
关闭线程池
shutdown():把状态从RUNNING设为SHUTDOWN,不再接收新任务,但 继续执行完队列中已有的任务,并中断空闲线程shutdownNow():把状态设为STOP,中断所有线程,并返回 队列中尚未执行的任务 列表
两者的区别在于 shutdown() 是「优雅关闭」,会跑完手头的工作;shutdownNow() 是「强制关闭」,立刻停止并返还未完成的任务。
Executors 工厂方法的坑
Executors 提供了几个快捷工厂方法,看起来方便,但在生产环境中往往有隐患:
// 固定线程数,无界队列 LinkedBlockingQueue
ExecutorService fixed = Executors.newFixedThreadPool(4);
// 核心线程数 0,最大线程数无上限,SynchronousQueue,空闲 60s 回收
ExecutorService cached = Executors.newCachedThreadPool();
// 单线程,无界队列
ExecutorService single = Executors.newSingleThreadExecutor();Executors.java
newFixedThreadPool/newSingleThreadExecutor使用 无界队列,任务堆积时可能耗尽内存(OOM),而且maximumPoolSize形同虚设newCachedThreadPool最大线程数为Integer.MAX_VALUE,任务量大时会创建海量线程,导致线程数失控甚至 OOM
因此规范做法是 手动 new ThreadPoolExecutor,显式传入有界队列和明确的拒绝策略,把资源和行为都控制在自己手里。
线程数如何设置
线程数量没有精确公式,但可以按任务的类型来估算:
| 任务类型 | 估算思路 |
|---|---|
| CPU 密集型 | 线程数 ≈ CPU 核数 + 1,多核并行计算 |
| IO 密集型 | 线程数 ≈ CPU 核数 * (1 + 等待时间 / 计算时间) |
- CPU 密集型任务(计算为主)开太多线程反而因为频繁上下文切换降低吞吐,通常用「核数 + 1」
- IO 密集型任务(等待网络/磁盘)大部分时间都在阻塞,可以开更多线程把等待时间利用起来
实际项目中建议 先用压测 + 监控(线程池活跃度、队列长度、拒绝次数)来验证,而不是拍脑袋定一个固定值。
Footnotes
-
还有一个只传
corePoolSize、maximumPoolSize、keepAliveTime、unit、workQueue的简化构造方法,内部会使用默认的DefaultThreadFactory和AbortPolicy。 ↩