Administrator
发布于 2019-02-21 / 5992 阅读
56

ThreadPoolExecutor 源码解析:execute 之后发生了什么

同事问了个我答不上来的问题

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.addt.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。同理 newCachedThreadPoolmaximumPoolSizeInteger.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-1order-2 之后,一眼就能看出谁在堆积。

另外加了监控,每分钟上报一次 getPoolSize()getQueue().size()getCompletedTaskCount(),队列长度超过 80% 就告警。

下篇预告

这篇先把《ThreadPoolExecutor 源码解析:execute 之后发生了什么》里的坑列了,下一篇写我们当时是怎么在线上工程里真正落地的——包括那次让领导拍桌的故障复盘。

参考