Skip to content

线程池源码级生命周期与线上排查

很多人学线程池只停在“七个参数”和“核心线程、队列、最大线程、拒绝策略”。这能应付一部分面试,但到了生产环境还不够。真正线上排查时,你会遇到:

  1. 明明线程池没满,接口还是慢。
  2. 队列一直涨,不知道是线程少还是下游慢。
  3. submit 的任务失败了,日志里却看不到异常。
  4. shutdown 后服务迟迟停不掉。
  5. jstack 里线程都是 WAITINGTIMED_WAITINGBLOCKED,不知道分别代表什么。
  6. 线程池越调越大,吞吐没上去,数据库反而被打满。

这篇把线程池从源码生命周期、任务流转、线程状态、监控指标、排查流程串成一条线。

学习目标

学完这一页,你要能回答:

  1. ThreadPoolExecutor 内部为什么同时管理“线程池状态”和“线程数量”。
  2. execute()addWorker()runWorker()getTask() 大致做什么。
  3. 一个任务从提交、入队、执行、异常、完成到线程复用的完整过程。
  4. 线程为什么会被回收,为什么有的线程一直不退出。
  5. executesubmit 的异常为什么表现不同。
  6. jstack 中不同线程状态如何对应线程池问题。
  7. 队列堆积时如何判断是线程池配置问题、下游瓶颈、锁等待、CPU 瓶颈还是提交速度过快。
  8. 线程池参数调整为什么必须联动数据库连接池、HTTP 连接池、Redis、MQ 和下游接口。

线程池不是一个简单队列

ThreadPoolExecutor 至少同时管理五件事:

mermaid
flowchart TD
    A["ThreadPoolExecutor"] --> B["线程池状态<br/>RUNNING / SHUTDOWN / STOP"]
    A --> C["工作线程数量<br/>workerCount"]
    A --> D["任务队列<br/>workQueue"]
    A --> E["Worker 集合<br/>workers"]
    A --> F["拒绝策略<br/>handler"]
    A --> G["线程工厂<br/>threadFactory"]

如果只把线程池理解成“几个线程加一个队列”,会漏掉两个关键点:

  1. 线程池要判断自己是否还能接收新任务。
  2. 线程池要控制工作线程数量不能超过边界。

这就是为什么源码里会把运行状态和线程数量一起管理。不同 JDK 实现细节可能有差异,但思想是一致的:状态决定能不能接任务,线程数量决定能不能创建 Worker

运行状态为什么重要

线程池常见状态:

状态是否接收新任务是否处理队列任务典型来源
RUNNING正常运行
SHUTDOWN调用 shutdown()
STOP否,并尝试中断执行中任务调用 shutdownNow()
TIDYING所有任务结束,准备终止
TERMINATED完全终止

流程图:

mermaid
flowchart TD
    A["RUNNING<br/>接收新任务并处理队列"] --> B["SHUTDOWN<br/>不接新任务但处理队列"]
    A --> C["STOP<br/>不接新任务并中断任务"]
    B --> D["TIDYING<br/>任务全部结束"]
    C --> D
    D --> E["TERMINATED<br/>线程池终止"]

为什么需要这些状态?

如果没有状态控制后果
关闭后仍然接收任务应用停机时任务越积越多
shutdown 直接杀任务已提交的订单、消息、文件处理可能丢失
shutdownNow 不尝试中断紧急停机时长时间卡住
不区分队列任务和新任务无法做到优雅停机

一个任务提交后的完整生命周期

execute() 为例,可以把流程拆成三段。

第一段:是否直接创建核心 Worker

mermaid
flowchart TD
    A["execute(command)"] --> B{"workerCount < corePoolSize"}
    B -- "是" --> C["addWorker(command, true)"]
    C --> D{"创建成功"}
    D -- "是" --> E["新 Worker 执行 command"]
    D -- "否" --> F["进入后续入队流程"]
    B -- "否" --> F

这里的重点是:只要当前工作线程数小于核心线程数,线程池倾向于创建新 Worker,而不是先排队。

第二段:核心线程满后尝试入队

mermaid
flowchart TD
    A["核心线程已满"] --> B{"线程池仍是 RUNNING"}
    B -- "否" --> C["拒绝任务"]
    B -- "是" --> D{"workQueue.offer(command)"}
    D -- "成功" --> E["任务进入队列"]
    E --> F["再次检查线程池状态"]
    F --> G{"如果关闭则移除并拒绝"}
    F --> H{"如果没有 Worker 则补一个空 Worker"}
    D -- "失败" --> I["尝试创建非核心 Worker"]

为什么入队后还要再次检查状态?

因为并发环境下,任务刚入队,另一个线程可能调用了 shutdown()。如果不二次检查,就可能出现线程池关闭后仍然留下新任务。

第三段:队列满后创建非核心 Worker 或拒绝

mermaid
flowchart TD
    A["队列已满"] --> B{"workerCount < maximumPoolSize"}
    B -- "是" --> C["addWorker(command, false)"]
    C --> D{"创建成功"}
    D -- "是" --> E["非核心 Worker 执行任务"]
    D -- "否" --> F["执行拒绝策略"]
    B -- "否" --> F

这就是“先核心线程、再队列、再最大线程、最后拒绝”的完整原因。

Worker 到底是什么

Worker 可以理解成线程池内部的工作单元。它不只是一个 Thread,还要记录第一个任务、执行状态,并参与线程池生命周期管理。

mermaid
flowchart TD
    A["Worker"] --> B["thread<br/>真正运行的线程"]
    A --> C["firstTask<br/>创建 Worker 时带入的第一个任务"]
    A --> D["completedTasks<br/>完成任务数"]
    A --> E["锁控制<br/>避免运行中被错误中断"]

为什么 Worker 需要 firstTask

因为创建线程的目的通常就是为了立刻执行当前提交的任务。如果每次都先创建空线程,再让线程去队列里取任务,会多一步入队和竞争。firstTask 可以让新 Worker 直接执行当前任务。

runWorker 循环:线程为什么能复用

线程池复用线程的核心在 runWorker 思想:一个 Worker 执行完第一个任务后,并不马上退出,而是继续从队列取任务。

mermaid
flowchart TD
    A["Worker 线程启动"] --> B["取 firstTask"]
    B --> C{"task 是否为空"}
    C -- "否" --> D["执行 task.run()"]
    D --> E["执行 afterExecute"]
    E --> F["task 置空"]
    F --> G["getTask 从队列取下一个任务"]
    G --> C
    C -- "是" --> H["Worker 退出"]
    H --> I["processWorkerExit"]

这就是为什么线程池不是“一个任务一个线程”。正确理解是:

一个 Worker 线程会执行多个任务,任务之间通过队列衔接。

如果任务里使用了 ThreadLocal 却不清理,下一个任务复用同一个线程时就可能读到旧值。

getTask:线程为什么会等待或退出

getTask() 的职责是从队列取任务,并决定当前 Worker 是否应该退出。

mermaid
flowchart TD
    A["getTask"] --> B{"线程池状态是否允许取队列任务"}
    B -- "不允许" --> C["减少 workerCount 并退出"]
    B -- "允许" --> D{"当前线程是否需要超时回收"}
    D -- "是" --> E["poll 等待 keepAliveTime"]
    D -- "否" --> F["take 一直等待任务"]
    E --> G{"是否拿到任务"}
    F --> G
    G -- "拿到" --> H["返回任务"]
    G -- "超时未拿到" --> I["判断是否退出 Worker"]

这里能解释几个现象:

现象原因
核心线程长期 WAITING默认核心线程不超时,队列为空时一直 take() 等任务
非核心线程一段时间后消失超过 keepAliveTime 没任务,被回收
开启 allowCoreThreadTimeOut(true) 后核心线程也会退出核心线程也使用超时等待
shutdown() 后线程池能处理完队列再退出状态不接新任务,但允许继续取队列任务

execute 和 submit 异常为什么不同

execute

execute() 提交的是 Runnable。任务抛出运行时异常时,异常会从工作线程里抛出,通常可以被线程的 UncaughtExceptionHandler 或日志看到,线程池会补充新的 Worker。

java
pool.execute(() -> {
    throw new RuntimeException("execute task failed");
});

submit

submit() 会把任务包装成 FutureTask。任务异常不会直接抛到工作线程外,而是保存到 FutureTask 内部,调用 future.get() 时才抛出 ExecutionException

java
Future<?> future = pool.submit(() -> {
    throw new RuntimeException("submit task failed");
});

future.get(); // 这里才能感知异常

对比:

提交方式异常去哪了风险
execute从任务线程抛出如果没设置线程异常处理器,日志可能不清晰
submit保存到 Future不调用 get() 就像“没失败”

商业项目里,异步入库、发送 MQ、生成报表这类任务如果用 submit 后不检查 Future,可能出现“任务失败但业务以为成功提交”的问题。

beforeExecute 和 afterExecute 有什么用

ThreadPoolExecutor 提供了钩子方法:

java
protected void beforeExecute(Thread t, Runnable r) { }

protected void afterExecute(Runnable r, Throwable t) { }

protected void terminated() { }

它们适合做:

  1. 统计任务耗时。
  2. 捕获和记录任务异常。
  3. 清理 ThreadLocal
  4. 上报线程池指标。
  5. 在线程池终止时释放资源。

Demo:

java
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;

public class TraceableThreadPool extends ThreadPoolExecutor {
    private static final ThreadLocal<Long> START_TIME = new ThreadLocal<Long>();

    public TraceableThreadPool() {
        super(
                4,
                8,
                60,
                TimeUnit.SECONDS,
                new ArrayBlockingQueue<Runnable>(200),
                r -> new Thread(r, "trace-worker-" + System.nanoTime()),
                new CallerRunsPolicy()
        );
    }

    @Override
    protected void beforeExecute(Thread thread, Runnable task) {
        START_TIME.set(System.currentTimeMillis());
    }

    @Override
    protected void afterExecute(Runnable task, Throwable throwable) {
        try {
            Long start = START_TIME.get();
            if (start != null) {
                long cost = System.currentTimeMillis() - start;
                System.out.println("task cost=" + cost + "ms");
            }
            if (throwable != null) {
                System.out.println("task error=" + throwable.getMessage());
            }
        } finally {
            START_TIME.remove();
        }
    }
}

注意:如果任务是 submit() 包装出来的 FutureTask,异常可能在 Future 里,afterExecutethrowable 不一定直接拿到。生产中要额外判断 Runnable 是否是 Future,再调用 get() 获取异常,但不要在普通未完成任务上阻塞等待。

线程状态和线程池问题怎么对应

jstack 看线程池时,不要只看线程数量,要看线程都卡在哪里。

线程状态常见含义线程池场景
RUNNABLE正在运行或等待 CPU/系统调用CPU 计算、网络读写、本地方法
WAITING无限期等待空闲 Worker 在队列 take() 等任务
TIMED_WAITING限时等待sleep、带超时的 poll、HTTP read timeout
BLOCKED等待进入 synchronized 锁多线程竞争同一把对象锁
WAITING on condition等条件满足等数据库连接、队列、Future 结果

场景一:大量线程 WAITING 在队列

可能是正常空闲:

text
java.lang.Thread.State: WAITING
at java.util.concurrent.locks.LockSupport.park
at java.util.concurrent.LinkedBlockingQueue.take

如果线程池没有任务,这是正常的。不要看到 WAITING 就以为线程死了。

场景二:大量线程卡在 socketRead

text
java.lang.Thread.State: RUNNABLE
at java.net.SocketInputStream.socketRead0

通常表示线程在等网络响应,例如 HTTP、数据库、Redis 或第三方接口。此时盲目增加线程池可能把下游打得更慢。

排查方向:

  1. 下游接口 P95/P99 是否升高。
  2. 是否设置连接超时和读取超时。
  3. HTTP 连接池是否打满。
  4. 是否有重试放大流量。

场景三:大量线程等待数据库连接

常见表现:

text
waiting on condition
at com.zaxxer.hikari.pool.HikariPool.getConnection

这说明线程池线程已经拿到执行机会,但卡在连接池。线程池调大不能解决数据库连接不足,可能造成更多线程排队。

处理方向:

  1. 查 SQL 是否慢。
  2. 查数据库连接池活跃数、等待数。
  3. 查事务是否过大或连接未关闭。
  4. 评估是否需要限流、拆批或优化索引。

场景四:大量线程 BLOCKED

说明线程在抢同一把 synchronized 锁。

处理方向:

  1. 找到锁对象和持有锁线程。
  2. 看锁内是否做了 IO、数据库、远程调用。
  3. 缩小锁范围。
  4. 用并发容器或更细粒度锁替代大锁。

线程池堆积排查决策树

mermaid
flowchart TD
    A["线程池任务堆积"] --> B["看提交 TPS 和完成 TPS"]
    B --> C{"提交 TPS 是否长期大于完成 TPS"}
    C -- "否" --> D["可能是瞬时峰值或监控窗口问题"]
    C -- "是" --> E["看 activeCount"]
    E --> F{"activeCount 是否接近 maximumPoolSize"}
    F -- "否" --> G["查任务是否没提交、调度是否异常、核心线程是否过小"]
    F -- "是" --> H["看线程栈"]
    H --> I{"线程主要卡在哪里"}
    I -- "CPU 计算" --> J["CPU 瓶颈,减少线程或优化算法"]
    I -- "数据库连接/慢 SQL" --> K["优化 SQL、连接池、事务、限流"]
    I -- "HTTP/Redis/下游接口" --> L["设置超时、熔断、隔离、降并发"]
    I -- "锁等待" --> M["缩小锁范围或改并发结构"]
    I -- "Future.get 等子任务" --> N["排查线程饥饿死锁"]
    K --> O["调整线程池和下游容量"]
    L --> O
    M --> O
    N --> O

关键判断:

指标说明
提交 TPS每秒进入线程池的任务数
完成 TPS每秒真正完成的任务数
队列长度积压任务数
activeCount正在执行任务的线程数
拒绝次数线程池过载次数
单任务耗时平均、P95、P99 耗时
下游耗时DB、Redis、HTTP、MQ 等耗时

如果提交 TPS 长期大于完成 TPS,队列必然越来越大。此时不是“调一个神奇参数”能解决,而是要提升完成 TPS、降低提交 TPS,或让系统明确拒绝/反压。

线程池参数调整的正确顺序

很多人一看到堆积就加 maximumPoolSize。这很危险。

正确顺序:

mermaid
flowchart TD
    A["线程池有问题"] --> B["先确认任务类型"]
    B --> C["看线程栈确认卡点"]
    C --> D["看下游容量"]
    D --> E["看队列等待时间"]
    E --> F["决定加线程、降并发、调队列或拒绝"]
    F --> G["压测验证"]
    G --> H["上线监控和回滚预案"]
现象不建议建议
DB 慢导致线程堆积直接加线程优化 SQL、限流、调连接池、拆批
HTTP 下游慢加大队列无限等超时、熔断、隔离、降级
CPU 100%增加线程减少线程、优化算法、扩容机器
队列过大用户超时继续加队列缩小队列、快速失败或反压
任务有依赖等待加大最大线程数拆线程池或改异步编排

商业场景:医疗采集任务堆积

假设采集平台每分钟拉取医院接口数据,任务进入采集线程池。

mermaid
flowchart TD
    A["调度中心触发采集"] --> B["提交到采集线程池"]
    B --> C["调用医院接口"]
    C --> D["清洗和校验"]
    D --> E["写 MySQL"]
    E --> F["发送 MQ 同步 ES"]
    F --> G["更新采集任务状态"]

如果线程池队列上涨,可能原因不是线程池本身:

卡点表现处理
医院接口慢线程卡在 socketRead设置超时、降并发、失败补偿
数据库慢线程等待连接或慢 SQL优化索引、批量写入、缩事务
ES/MQ 慢发送耗时升高异步化、失败表、重试和死信
清洗 CPU 高CPU 接近 100%控制线程数、优化规则、扩容
子任务等待Future.get 卡住拆池或 CompletableFuture 编排

可运行简化 Demo:

java
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;

public class CollectPoolDemo {
    private static final ThreadPoolExecutor COLLECT_POOL =
            new ThreadPoolExecutor(
                    8,
                    16,
                    60,
                    TimeUnit.SECONDS,
                    new ArrayBlockingQueue<Runnable>(500),
                    r -> new Thread(r, "collect-worker-" + System.nanoTime()),
                    new ThreadPoolExecutor.CallerRunsPolicy()
            );

    public static void main(String[] args) {
        for (int i = 0; i < 100; i++) {
            int taskId = i;
            COLLECT_POOL.execute(() -> collect(taskId));
        }
        COLLECT_POOL.shutdown();
    }

    private static void collect(int taskId) {
        try {
            mockHospitalApi();
            mockSaveDb();
            System.out.println(Thread.currentThread().getName() + " finish task " + taskId);
        } catch (Exception e) {
            System.out.println("collect failed taskId=" + taskId + ", error=" + e.getMessage());
        }
    }

    private static void mockHospitalApi() throws InterruptedException {
        TimeUnit.MILLISECONDS.sleep(100);
    }

    private static void mockSaveDb() throws InterruptedException {
        TimeUnit.MILLISECONDS.sleep(50);
    }
}

这个 Demo 不是为了模拟真实医院接口,而是体现生产原则:

  1. 线程池有界。
  2. 线程有业务名称。
  3. 满载时 CallerRunsPolicy 反压提交方。
  4. 任务内部要捕获并记录失败。
  5. 真实项目还要加超时、幂等、失败表和告警。

面试标准回答

线程池任务堆积怎么排查

标准回答:

text
线程池堆积不能只看线程数。先看提交 TPS、完成 TPS、队列长度、activeCount、拒绝次数和单任务耗时。如果提交速度长期大于完成速度,队列一定上涨。然后用 jstack 看工作线程卡在哪里:如果卡在 socketRead,查下游接口和超时;卡在数据库连接,查连接池、慢 SQL 和事务;大量 BLOCKED 查锁竞争;大量 Future.get 查线程饥饿死锁;CPU 高则查计算和上下文切换。处理时要结合下游容量,不能盲目加线程,否则可能把数据库或第三方接口打得更慢。

追问点:

  1. activeCount 满但 CPU 不高说明什么?
  2. 队列很大但没有拒绝是不是好事?
  3. 为什么线程池调大后吞吐不升反降?
  4. 如何判断是线程池瓶颈还是数据库瓶颈?

线程池源码执行链路怎么讲

标准回答:

text
ThreadPoolExecutor 提交任务后,先判断当前 workerCount 是否小于 corePoolSize,小于就创建 Worker 执行任务;否则尝试把任务放入 workQueue;入队成功后还要二次检查线程池状态,防止关闭后残留新任务;如果队列满了,再判断是否能创建到 maximumPoolSize;如果不能就执行拒绝策略。Worker 启动后会先执行 firstTask,然后在 runWorker 循环中不断通过 getTask 从队列取任务。getTask 会根据线程池状态、keepAliveTime、allowCoreThreadTimeOut 判断线程是等待任务还是退出回收。

追问点:

  1. 为什么无界队列会让最大线程数基本失效?
  2. getTask() 为什么有时 take,有时 poll
  3. 为什么入队后还要二次检查线程池状态?
  4. submit 为什么可能隐藏异常?

关联知识点

知识点说明
线程池总览线程池学习路线
核心参数七大参数和配置后果
执行流程原理execute 主流程
队列与拒绝策略过载时如何排队、反压、拒绝
参数估算与监控怎么按业务容量配置线程池
关闭与常见坑shutdown、submit 异常、饥饿死锁
ThreadLocal线程复用下为什么必须 remove
JVM 排查工具jstack、jmap、jstat、dump 分析

本章小结

线程池是一个容量控制系统,不是一个简单异步工具。真正掌握线程池,要能把源码流程和生产现象对应起来:execute 决定任务怎么进入系统,Worker 决定线程怎么复用,getTask 决定线程怎么等待和退出,队列和拒绝策略决定过载时系统怎么保护自己,jstack 和监控指标决定你能不能在事故中找到真正瓶颈。