跳转到内容
唯一赫兹
返回

ThreadPoolExecutor 线程池:核心参数与执行流程

为什么需要线程池

如果每个任务都 new Thread() 去执行:创建和销毁线程本身就有开销,而且销毁后线程无法复用,CPU 缓存也就失效了。线程池的核心思路是 复用线程:预先创建好一批线程,任务来了就交给空闲线程去跑,避免反复创建销毁线程。

判断一个优秀线程池的标准:

  1. 复用:任务交给已有的空闲线程,而不是每次都新建
  2. 控制并发数:避免无限制创建线程耗尽资源
  3. 管理生命周期:空闲线程可以被回收,池可以被优雅关闭

核心参数

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超过核心线程数部分的线程,空闲超过该时间就会被回收
unitkeepAliveTime 的时间单位
workQueue任务队列,当线程数达到 corePoolSize 后,新任务先入队
threadFactory线程工厂,用于创建线程,可自定义线程名/优先级/是否守护线程
handler拒绝策略,线程池无法接收新任务时的处理方式

任务来了怎么走:如果线程数少于 corePoolSize,直接创建新线程执行;如果已经达到 corePoolSize,任务放入 workQueue 等待;如果队列满了,再创建线程直到 maximumPoolSize;如果线程数已经达到 maximumPoolSize 且队列也满了,就交给拒绝策略处理。

工作队列 workQueue

workQueue 决定了「任务排队」的规则,它直接影响了线程池如何增长线程。常用的队列有:

队列类型是否有界特点
ArrayBlockingQueue基于数组有界(初始化指定)适合对任务堆积数有明确上限、资源可控的场景
LinkedBlockingQueue基于链表默认无界,可传容量可以无界;newFixedThreadPool 默认用它
SynchronousQueue移交队列容量为 0本身不存任务,每次 put 必须等一个 take 配对,用于直接移交
PriorityBlockingQueue优先级队列默认可扩容任务按优先级出队
DelayQueue延迟队列无界任务到期才能出队,常用于调度

队列类型如何影响线程增长? 上一节的执行流程里有一个隐含前提:队列是否容易满,决定了线程数会不会往 maximumPoolSize 涨。

所以「有界队列 + 明确的拒绝策略」通常是生产环境的稳妥组合:它既能限制堆积量、防止 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

把状态和数量压缩进一个 int,就能用一次 CAS 同时更新「状态」和「线程数」,保证两者的一致性和原子性。

线程池的五种状态

状态之间的大小关系是 RUNNING < SHUTDOWN < STOP < TIDYING < TERMINATED,所以源码里用 runStateAtLeast(c, STOP) 这种数值比较来判断状态:

状态含义
RUNNING正常运行:接收新任务,处理队列里的任务
SHUTDOWN调用 shutdown() 后:不接收新任务,但会处理队列里已有的任务
STOP调用 shutdownNow() 后:不接收新任务,不处理队列任务,并中断正在执行的任务
TIDYING所有任务已结束、工作线程数为 0,即将执行 terminated() 钩子方法
TERMINATEDterminated() 已执行完毕,线程池彻底终止

SHUTDOWNSTOP 的关键区别在于: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

整个逻辑可以概括为三步:核心线程 → 入队 → 非核心线程 → 拒绝

其中第二步的「二次检查」很关键:任务入队成功不代表万事大吉。入队之后到检查之前,线程池可能已经被其他线程关闭了,此时必须把任务从队列里移除再交给拒绝策略;同时还要防止因核心线程为空(比如被超时回收或正在退出)而出现「队列里有任务却没人跑」的情况——这时要补一个线程,它的 firstTasknull,因为它没有任何初始任务,直接进入 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 分为两段:

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 是否正在执行任务」

这样就能做到「只中断空闲线程,不打断正在跑任务的线程」。

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

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 的作用在这里体现:

所以「超过核心线程数的那些线程在空闲 keepAliveTime 后会被回收」。如果设置了 allowCoreThreadTimeOut = true,核心线程也会被回收。

拒绝策略

当任务无法提交(线程数达到 maximumPoolSize 且队列已满,或线程池已关闭)时,会交给 RejectedExecutionHandler 处理。JDK 内置了四种策略:

策略行为风险
AbortPolicy默认策略,直接抛出 RejectedExecutionException调用方需自行捕获异常
CallerRunsPolicy提交任务的线程 自己执行该任务提交线程被占用,可能拖慢流程
DiscardPolicy静默丢弃新任务任务可能丢失,无任何提示
DiscardOldestPolicy丢弃队列中 最老 的任务,然后重试提交新任务可能丢任务,适合时效性高的场景

CallerRunsPolicy 是一种「限流」手段:当线程池跑不动时,让提交者自己跑,一方面起到背压作用,另一方面任务不会被丢弃。

关闭线程池

两者的区别在于 shutdown() 是「优雅关闭」,会跑完手头的工作;shutdownNow() 是「强制关闭」,立刻停止并返还未完成的任务。

Executors 工厂方法的坑

Executors 提供了几个快捷工厂方法,看起来方便,但在生产环境中往往有隐患:

// 固定线程数,无界队列 LinkedBlockingQueue
ExecutorService fixed = Executors.newFixedThreadPool(4);

// 核心线程数 0,最大线程数无上限,SynchronousQueue,空闲 60s 回收
ExecutorService cached = Executors.newCachedThreadPool();

// 单线程,无界队列
ExecutorService single = Executors.newSingleThreadExecutor();Executors.java

因此规范做法是 手动 new ThreadPoolExecutor,显式传入有界队列和明确的拒绝策略,把资源和行为都控制在自己手里。

线程数如何设置

线程数量没有精确公式,但可以按任务的类型来估算:

任务类型估算思路
CPU 密集型线程数 ≈ CPU 核数 + 1,多核并行计算
IO 密集型线程数 ≈ CPU 核数 * (1 + 等待时间 / 计算时间)

实际项目中建议 先用压测 + 监控(线程池活跃度、队列长度、拒绝次数)来验证,而不是拍脑袋定一个固定值。

Footnotes

  1. 还有一个只传 corePoolSizemaximumPoolSizekeepAliveTimeunitworkQueue 的简化构造方法,内部会使用默认的 DefaultThreadFactoryAbortPolicy



上一篇
如何使用 AI Agent 管理大型项目
下一篇
Java 中断机制