Skip to content
Go back

Java线程池——execute()为什么入队后要做二次检查?

线程池: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 位、ctlOfrs | 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),它从队列中拉取任务消费。这里传 nullcore=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 有独特的”倒流”效果:任务提交速度太快 → 提交线程自己执行任务 → 提交速度自然降低 → 天然背压效果。DiscardOldestPolicye.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(),无需计算线程数。

章末提问


Share this post on:

Previous Post
LongAdder vs AtomicLong——高并发下的Cell分槽与伪共享
Next Post
ForkJoinPool分治并行:工作窃取与双端队列