JAVA · Vol.II · DAY 15 · 并发编程
线程池 ThreadPoolExecutor:7 大参数与 4 种拒绝策略
凌晨 2 点的告警:200 个请求同时超时
"线上服务 200 个请求同时超时,jstack 打出来一看——线程池队列满了,200 个线程全部 WAITING。" 这不是段子,是真实 P0 事故。
那天晚上,运营做了一波大促预热,流量比平时翻了 5 倍。我们的订单查询服务使用了一个"看起来挺合理"的线程池配置:
// 当时的"合理"配置 ExecutorService pool = new ThreadPoolExecutor( 10, // corePoolSize 10, // maximumPoolSize —— 和 core 一样! 0L, TimeUnit.SECONDS, new LinkedBlockingQueue<>() // 无界队列!默认容量 Integer.MAX_VALUE );
问题一目了然:LinkedBlockingQueue 默认容量是 Integer.MAX_VALUE(约 21 亿),队列永远不会满,所以 maximumPoolSize 形同虚设——线程池永远只有 10 个核心线程在干活。当流量洪峰涌入,10 个线程处理不过来,任务全部堆在队列里排队,每个请求都在等队列里的任务被执行,而队列里的任务在等线程来取——死等,超时,雪崩。
更隐蔽的一层:这个配置里 maximumPoolSize == corePoolSize,等于亲手把「临时工」这条路焊死了。线程池的弹性完全来自「队列满 → 创建非核心线程」这一步,而这条触发路径的前提是队列必须有可能满。无界队列让 maximumPoolSize 永远轮不到上场。
线程池是并发编程面试出现频率最高的主题,没有之一。原因有三:
- 覆盖面广——7 个参数串联了线程管理、队列、锁、OOM、GC 等核心知识点
- 生产强相关——几乎每个 Java 后端项目都在用,配置错误真的会炸
- 区分度高——能说清楚 execute() 决策树的候选人,基本功不会差
7 大参数详解:把线程池想象成一家餐厅
ThreadPoolExecutor 的构造函数有 7 个参数,死记硬背容易混。我们用一个类比来串起来:
| 参数 | 餐厅类比 | 技术含义 | 默认值/常见值 |
|---|---|---|---|
corePoolSize | 正式厨师 | 核心线程数,即使空闲也不会被回收 | 按业务设定 |
maximumPoolSize | 正式 + 临时工上限 | 池中允许的最大线程数 | ≥ corePoolSize |
keepAliveTime | 临时工空闲多久下班 | 非核心线程空闲存活时间 | 60s |
unit | 时间的单位 | keepAliveTime 的时间单位 | SECONDS |
workQueue | 等候区 | 存放待执行任务的阻塞队列 | LinkedBlockingQueue |
threadFactory | 招聘渠道 | 创建线程的工厂,常用于设置线程名 | Executors.defaultThreadFactory() |
rejectedExecutionHandler | 客满怎么处理 | 线程池和队列都满时的拒绝策略 | AbortPolicy |
corePoolSize 是"保底编制",maximumPoolSize 是"最大编制",keepAliveTime 只管非核心线程的去留,workQueue 是核心和最大线程之间的缓冲区。记住决策顺序:核心线程 → 队列 → 非核心线程 → 拒绝,这个顺序比每个参数的定义更重要。
execute() 决策树:一个任务的生死之旅
当你调用 pool.execute(task) 时,线程池内部会经过一系列判断来决定这个任务的命运。这段逻辑是整个 ThreadPoolExecutor 的灵魂,面试必问。
public void execute(Runnable command) { if (command == null) throw new NullPointerException(); int c = ctl.get(); // ctl 是 AtomicInteger,高3位=状态,低29位=线程数 int workerCount = workerCountOf(c); // 第一步:线程数 < corePoolSize → 创建核心线程 if (workerCount < corePoolSize) { if (addWorker(command, true)) // 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); // 兜底:确保至少有一个线程消费队列 return; } // 第三步:队列满 → 创建非核心线程 if (!addWorker(command, false)) // false = 创建非核心线程 reject(command); // 第四步:线程也满了 → 拒绝 }
面试官追问:为什么入队成功后还要 recheck 一次线程池状态?
因为在 offer() 成功和 recheck 之间,线程池可能被调用了 shutdown()。如果此时线程池已关闭且任务可以从队列中移除,就拒绝该任务。如果线程池还在运行但 workerCount 已经为 0(所有线程都意外退出了),就创建一个非核心线程来保证队列中的任务能被消费——这是防止「队列里有任务但没人消费」的死锁式僵局。
追问二:为什么 addWorker 失败后要重新 get 一次 ctl?
因为 addWorker 内部要用 CAS 改 ctl(线程数 +1),并发下可能失败(别的线程同时加了线程)。失败后线程数可能已经变了,必须重新读 ctl 再做后续判断,否则会用过期数据走错分支。这也是 ctl 被设计成「状态 + 线程数」打包在一个 AtomicInteger 里的原因之一——状态变化和线程数变化可以在同一次原子操作里完成,下一站细讲。
workQueue 选择:5 种队列的取舍
workQueue 的选择直接决定了线程池在流量高峰期的表现。这是很多候选人容易忽略的点——他们能背出 7 个参数名,却说不出不同队列的适用场景。
| 队列类型 | 容量 | 特点 | 风险 | 适用场景 |
|---|---|---|---|---|
LinkedBlockingQueue |
默认 Integer.MAX_VALUE(≈无界) | 链表实现,put/take 两把锁分离 | 无界时 OOM | 必须指定容量后使用 |
ArrayBlockingQueue |
有界,必须指定 capacity | 数组实现,一把 ReentrantLock | 队列满时触发 max 线程或拒绝 | 生产环境首选 |
SynchronousQueue |
0(不存储任务) | 直接交付,必须「一单一线程」 | 没有空闲线程就直接走拒绝策略 | Executors.newCachedThreadPool() |
PriorityBlockingQueue |
无界 | 按 Comparable 优先级排序 | 低优先级任务可能饥饿 | 任务有明确优先级(VIP 订单) |
DelayQueue |
无界 | 元素实现 Delayed,到期才能 take |
无界同样要当心 OOM | 延迟任务:订单超时关单、定时重试 |
// 推荐:有界队列 + CallerRunsPolicy = 天然背压 new ThreadPoolExecutor( 20, // core 50, // max 60L, TimeUnit.SECONDS, // 非核心线程空闲 60s 回收 new ArrayBlockingQueue<>(1000), // 有界队列,容量 1000 new ThreadFactoryBuilder() // Guava 工具 .setNameFormat("order-pool-%d") .setDaemon(true) .build(), new ThreadPoolExecutor.CallerRunsPolicy() // 队列满 + 线程满 → 调用者线程自己跑 );
为什么 ArrayBlockingQueue 用一把锁而 LinkedBlockingQueue 用两把锁?
LinkedBlockingQueue 内部使用 putLock 和 takeLock 两把锁分离生产和消费操作,在高并发下吞吐量更高。ArrayBlockingQueue 只有一把 lock,但实现更简单、内存占用更低,且强制有界——对于线程池场景,这点吞吐量差异远不如"有界"带来的安全性重要。
生产环境默认 ArrayBlockingQueue + 合理容量;要削峰填谷且能接受排队就用它;要「来一个任务立刻分配线程」的纯并发场景才用 SynchronousQueue(配合大 max);有优先级或延迟语义再考虑另外两种。任何无界队列都要先问一句:内存扛得住吗?
4 种拒绝策略:任务被拒后的 4 种命运
当线程池的线程数已达 maximumPoolSize 且 workQueue 已满,新提交的任务将被拒绝。ThreadPoolExecutor 内置了 4 种拒绝策略:
// 1. AbortPolicy(默认)—— 直接抛异常 public static class AbortPolicy implements RejectedExecutionHandler { public void rejectedExecution(Runnable r, ThreadPoolExecutor e) { throw new RejectedExecutionException( "Task " + r.toString() + " rejected from " + e.toString()); } } // 2. CallerRunsPolicy —— 谁提交谁执行(反压利器) public static class CallerRunsPolicy implements RejectedExecutionHandler { public void rejectedExecution(Runnable r, ThreadPoolExecutor e) { if (!e.isShutdown()) { r.run(); // 在调用者线程直接执行,不抛异常,不丢任务 } } } // 3. DiscardPolicy —— 默默丢弃,不声张 public static class DiscardPolicy implements RejectedExecutionHandler { public void rejectedExecution(Runnable r, ThreadPoolExecutor e) { // 空实现,任务被静默丢弃 } } // 4. DiscardOldestPolicy —— 丢弃队列头部最旧的任务 public static class DiscardOldestPolicy implements RejectedExecutionHandler { public void rejectedExecution(Runnable r, ThreadPoolExecutor e) { if (!e.isShutdown()) { e.getQueue().poll(); // 丢弃队列头部的任务 e.execute(r); // 重新尝试提交当前任务 } } }
| 策略 | 行为 | 适用场景 | 风险 |
|---|---|---|---|
| AbortPolicy | 抛 RejectedExecutionException | 快速失败,让调用方感知 | 上层必须 catch,否则请求 500 |
| CallerRunsPolicy | 调用者线程自己执行任务 | Web 服务反压:自动降低提交速度 | 调用者线程(如 Tomcat 线程)被阻塞 |
| DiscardPolicy | 静默丢弃 | 日志采集、监控上报等可丢失场景 | 任务丢失无任何通知,排查困难 |
| DiscardOldestPolicy | 丢弃最旧任务,重试提交 | 只关心最新数据的场景(如股价推送) | 旧任务丢失,不适合要求严格顺序的业务 |
假设 Tomcat 工作线程向线程池提交任务,队列满了。此时 CallerRunsPolicy 会让 Tomcat 线程自己执行这个任务——相当于 Tomcat 线程被"征用"了,在它执行完之前,它无法处理下一个 HTTP 请求。这就自然地减缓了任务提交的速度,形成了一种负反馈机制:下游处理不过来 → 上游自动减速。这比直接抛异常然后返回 500 优雅得多。这也是「有界队列 + CallerRuns」成为生产默认组合的原因。
"生产环境一般用 CallerRunsPolicy 做反压。如果任务不能丢且不能阻塞调用者,我会自定义 RejectedExecutionHandler,把被拒绝的任务持久化到数据库或 MQ,事后补偿——具体写法在后面的站点给出。"
线程池的 5 种状态与 ctl 位运算
前面反复出现的 ctl 到底是什么?「线程池状态」又是哪来的?其实 ThreadPoolExecutor 用一个 int 就同时存了「池子状态」和「线程数」两件事。
线程池和线程一样也有生命周期:正在运行 → 已关闭(只处理存量)→ 停止(中断存量)→ 收尾 → 终结。JDK 用 5 个状态描述它:
这 5 个状态就存在 ctl 这个 volatile AtomicInteger 的高 3 位里,低 29 位存线程数:
private final AtomicInteger ctl = new AtomicInteger(ctlOf(RUNNING, 0)); private static final int COUNT_BITS = 9; // 2^29 = 5 亿,线程数上限 private static final int CAPACITY = (1 << COUNT_BITS) - 1; // 高 3 位的状态值,按数值大小排好序,isRunning(c) 这类判断才成立 private static final int RUNNING = -1 << COUNT_BITS; // 正常运行,接受新任务 private static final int SHUTDOWN = 0 << COUNT_BITS; // 已 shutdown,不接新任务,处理存量 private static final int STOP = 1 << COUNT_BITS; // 已 stop,中断运行中任务,丢弃队列 private static final int TIDYING = 2 << COUNT_BITS; // 所有线程已退出,workerCount == 0 private static final int TERMINATED = 3 << COUNT_BITS; // terminated() 钩子已执行 private static int workerCountOf(int c) { return c & CAPACITY; } // 取低 29 位 = 线程数 private static boolean isRunning(int c) { return c < SHUTDOWN; } // 高 3 位 < 000 只能是 111
为什么要打包成一个 int?两个原因:
- 原子性:
addWorker里「状态检查 + 线程数 +1」必须是一个原子动作。CAS 一次搞定,避免「状态还是 RUNNING,但线程数已经超了」的中间态。 - 可见性:
volatile语义保证任何一个线程改了ctl,其他线程立刻看得到——shutdown 的传播就靠它。
「一个 int 存两件事」是 AQS 之外的又一经典位运算设计(synchronized 的对象头 Mark Word 也是同理)。面试答到 ctl 时,先说高 3 位 5 状态、低 29 位线程数,再说「为什么打包」——区分度立刻出来。
线程池里的异常去哪儿了:execute vs submit
任务在池子里抛了异常,会发生什么?两种提交方式的答案完全不同——这也是线上「任务悄悄失败、半天没人发现」的头号原因。
// ① execute:异常从工作线程"飞"出来 pool.execute(() -> { int x = 1 / 0; // ArithmeticException }); // 提交线程完全无感。异常处理路径: // 工作线程死亡 → UncaughtExceptionHandler 打印堆栈(默认打到 stderr) // → 线程池再创建一个新工作线程替补。 // 如果没人处理 UncaughtExceptionHandler,异常就"蒸发"了——不告警、不记录, // 任务失败这件事只有翻日志才知道。 // ② submit:异常被"吞"进 Future Future<String> f = pool.submit(() -> { int x = 1 / 0; return "ok"; }); // f.isDone() == true,看起来任务"正常完成"了! // 异常被包装成 ExecutionException 存在 Future 里, // 只有调用 f.get() 才会抛出来。 // 如果业务代码从不 get()(比如 fire-and-forget 场景), // 这个异常就永远不可见。
| execute(Runnable) | submit(Callable/Runnable) | |
|---|---|---|
| 异常去向 | 抛到 UncaughtExceptionHandler | 存进 Future,等 get() |
| 提交方感知 | 无(异步) | get() 时抛 ExecutionException |
| 不处理后果 | 静默失败 + 工作线程死亡重建 | 静默失败,Future 永远不返回结果 |
| 适用 | 不关心结果的 fire-and-forget | 需要结果或需要异常感知的任务 |
// 兜底 1:任务体内永远 try-catch(最可靠,两种提交都管用) pool.execute(() -> { try { doBusiness(); } catch (Exception e) { log.error("订单任务执行失败, orderId={}", orderId, e); metrics.inc("task.failed"); // 打点告警 } }); // 兜底 2:ThreadFactory 挂 UncaughtExceptionHandler(抓 execute 的漏网之鱼) new ThreadFactoryBuilder() .setNameFormat("order-pool-%d") .setUncaughtExceptionHandler((t, e) -> log.error("线程 {} 未捕获异常", t.getName(), e)) .build() // 兜底 3:submit 的结果要么 get/join,要么 exceptionally 处理 Future<String> f = pool.submit(task); f.get(5, TimeUnit.SECONDS); // 带超时,顺便防住任务卡死
我们团队出过「对账任务每天悄悄失败两天」的事故:submit 提交的对账任务里 NPE 了,没人 get,没人告警,直到财务发现差额。之后定下规矩:池内任务必须 try-catch + 打点,提交后必须有人消费结果。线程池不会帮你发现业务异常,它只负责把异常送到「你以为有人接、其实没人接」的地方。
ThreadFactory:线程命名是你排查问题的第一张牌
第 7 站的代码里用了 Guava 的 ThreadFactoryBuilder,但很多人默认「不传 threadFactory 也行」。不行——默认工厂创建的线程叫 pool-1-thread-1、pool-2-thread-3,一个服务里十个池子,jstack 打出来你根本分不清谁是谁。
public class NamedThreadFactory implements ThreadFactory { private final AtomicInteger seq = new AtomicInteger(1); private final String prefix; public NamedThreadFactory(String prefix) { this.prefix = prefix; } @Override public Thread newThread(Runnable r) { Thread t = new Thread(r, prefix + "-" + seq.getAndIncrement()); t.setDaemon(false); // 业务池建议非守护:JVM 退出前把任务跑完 return t; } }
线程名规范建议(我们团队的约定):
- 业务域-用途-序号:
order-pay-pool-1、order-notify-pool-2——看到名字就知道「哪个业务、干什么、第几个线程」。 - 序号用
AtomicInteger保证唯一,别用System.currentTimeMillis()这种。 - daemon 要想清楚:业务线程池一般设非守护(否则 JVM 可能在任务没跑完时就退出);纯后台辅助线程(如指标上报)可设守护。
jstack 抓到线程 dump 后,第一眼看线程名:一片 http-nio-8080-exec-* 都 BLOCKED 在同一个 monitor → 锁竞争;某个 order-pay-pool-* 集体 WAITING 在 SocketRead → 下游 RPC 卡了。命名清晰的池子让这类判断从「考古」变成「一眼」。
allowCoreThreadTimeOut 与线程预热
默认规则里有个容易忽略的细节:核心线程空闲也不回收。哪怕池子里 40 个核心线程一整天没活干,它们也一直占着内存等着。只有非核心线程才受 keepAliveTime 管辖。
// 默认:核心线程永生 pool.allowCoreThreadTimeOut(true); // 之后核心线程空闲超过 keepAliveTime 也会被回收。 // 适用:流量波动剧烈的服务——白天 40 个线程忙死,凌晨 1 个就够, // 回收核心线程能省下几十 MB 栈内存(每线程默认 1MB 栈)。 // 代价:流量突然回升时要重新创建线程(毫秒级,通常可接受)。
和它配对的还有「线程预热」:核心线程默认是懒创建的——池子刚 new 出来时一个线程都没有,前几个任务是边创建线程边执行的,第一次响应会偏慢。
// 单个预热:创建一个空闲核心线程 pool.prestartCoreThread(); // 全部预热:把 corePoolSize 个核心线程一次性建好 int started = pool.prestartAllCoreThreads(); // 返回本次真正新建的线程数。放在 @PostConstruct / 应用启动钩子里, // 避免第一个用户为你的线程池付「建线程」的成本。
常规服务:core 固定不回收 + 启动预热,行为最可预测。潮汐型服务(白天忙、夜里闲):开 allowCoreThreadTimeOut(true),把 keepAliveTime 设 30~60s,夜里自动缩到最小。这两个开关都不影响功能,只影响资源曲线和首包延迟。
Executors 工厂方法的陷阱:为什么阿里规约禁止使用
JDK 的 Executors 工具类提供了几个便捷的工厂方法,但它们每一个都藏着生产级的坑。
// 陷阱 1:newFixedThreadPool —— 无界队列 → OOM public static ExecutorService newFixedThreadPool(int nThreads) { return new ThreadPoolExecutor( nThreads, nThreads, 0L, TimeUnit.MILLISECONDS, new LinkedBlockingQueue<>() // ← 无界!任务堆积 → OOM ); } // 陷阱 2:newSingleThreadExecutor —— 同样的无界队列 public static ExecutorService newSingleThreadExecutor() { return new ThreadPoolExecutor( 1, 1, 0L, TimeUnit.MILLISECONDS, new LinkedBlockingQueue<>() // ← 无界! ); } // 陷阱 3:newCachedThreadPool —— 无限线程 → 线程爆炸 public static ExecutorService newCachedThreadPool() { return new ThreadPoolExecutor( 0, Integer.MAX_VALUE, // ← 最大线程数 21 亿! 60L, TimeUnit.SECONDS, new SynchronousQueue<>() // 没有缓冲,每个任务必须分配线程 ); }
| 工厂方法 | 核心问题 | 后果 | 替代方案 |
|---|---|---|---|
newFixedThreadPool |
LinkedBlockingQueue 无界 | 任务堆积 → OutOfMemoryError | 手动 new + ArrayBlockingQueue |
newSingleThreadExecutor |
LinkedBlockingQueue 无界 | 同上 | 手动 new,core=1 + 有界队列 |
newCachedThreadPool |
maxPoolSize = Integer.MAX_VALUE | 流量尖刺 → 创建大量线程 → OOM 或 CPU 100% | 手动 new + 合理 maxPoolSize |
【强制】线程池不允许使用 Executors 去创建,而是通过 ThreadPoolExecutor 的方式。这样的处理方式让写的同学更加明确线程池的运行规则,规避资源耗尽的风险。
真实踩坑:newFixedThreadPool 导致 Full GC 不停
某团队用 Executors.newFixedThreadPool(20) 做异步日志写入。上线初期一切正常,直到某天日志量大增,队列堆积了 200 万个任务对象,每个任务持有请求上下文约 2KB,总计占用约 4GB 堆内存。JVM 进入 Full GC 循环,服务假死。换成 ArrayBlockingQueue(5000) + CallerRunsPolicy 后问题解决。
参数调优:公式给起点,压测给答案
面试中被问"线程池参数怎么设",回答"看情况"是不够的。你需要给出公式、给出数字、给出验证手段。
CPU 密集型(纯计算、加密、序列化):
corePoolSize = N_CPU + 1
多一个线程是为了某个线程偶尔因缺页中断等暂停时,额外的线程能顶上。
IO 密集型(RPC 调用、数据库查询、HTTP 请求):
corePoolSize = N_CPU × (1 + W / C)
其中 W = 线程等待时间(等 IO),C = 线程计算时间(CPU 运算)。
/* * 场景:订单查询服务,每次请求要调 3 个 RPC + 2 次 DB 查询 * 机器:8 核 * 观测:单次请求中 CPU 计算约 10ms,IO 等待约 40ms * → W/C = 40/10 = 4 * * 套用公式:corePoolSize = 8 × (1 + 4) = 40 */ int cpuCores = Runtime.getRuntime().availableProcessors(); // 8 double wcRatio = 4.0; // 等待时间 / 计算时间 int coreSize = (int) (cpuCores * (1 + wcRatio)); // 40 int maxSize = coreSize * 2; // 80,留弹性空间 ThreadPoolExecutor pool = new ThreadPoolExecutor( coreSize, maxSize, 60L, TimeUnit.SECONDS, new ArrayBlockingQueue<>(2000), new ThreadFactoryBuilder().setNameFormat("order-q-%d").build(), new ThreadPoolExecutor.CallerRunsPolicy() );
真实调优经历:公式算出 40,实际只用了 25
某服务按公式设了 corePoolSize=40,但上线后发现线程切换开销明显——CPU 利用率只有 60% 但 load average 很高。原因:下游 RPC 超时设得很短(200ms),大多数 IO 等待时间其实没有 40ms 那么长。调低到 25 后,吞吐量反而提升了 15%。结论:公式给起点,监控给方向,压测给答案。W/C 别靠猜,用 APM 或日志把真实等待时间量出来。
监控与告警:给线程池装上仪表盘
公式只是起点,真正上线后必须配合监控数据持续调优。ThreadPoolExecutor 自带一组观测方法,全部真实存在:
// 定时采集,上报到 Prometheus / Grafana ScheduledExecutorService monitor = ...; monitor.scheduleAtFixedRate(() -> { log.info("[pool-monitor] active={}, poolSize={}, largest={}, queue={}, completed={}", pool.getActiveCount(), // 正在执行任务的线程数 pool.getPoolSize(), // 当前池中线程总数 pool.getLargestPoolSize(), // 历史峰值线程数(评估 max 设得合不合理) pool.getQueue().size(), // 队列中待执行的任务数 pool.getCompletedTaskCount() // 已完成的任务总数 ); }, 0, 10, TimeUnit.SECONDS);
| 监控指标 | 告警阈值 | 说明 |
|---|---|---|
getActiveCount() / getPoolSize() | > 80% | 线程利用率过高,考虑扩容 |
getQueue().size() / capacity | > 70% | 队列堆积,有拒绝风险 |
| 拒绝次数(自定义 handler 计数,见下一站) | > 0 | 出现拒绝,必须立即处理 |
getCompletedTaskCount() 增量 | 突降 | 吞吐量下降,可能有死锁或下游故障 |
getLargestPoolSize() | 长期贴近 max | max 偏小,或上游流量该治理了 |
为什么表里没有 getRejectedExecutionCount()?
因为JDK 压根没有这个方法——网上不少博客会引用它,是抄来抄去造出来的。线程池拒绝了多少任务,JDK 不帮你记,正确姿势是自定义 RejectedExecutionHandler 自己计数(下一站给出完整代码)。记住这个点,面试被追问「你怎么监控拒绝次数」时不会翻车。
自定义拒绝策略:计数 + 兜底落盘
四种内置策略覆盖不了「任务不能丢、又不能阻塞调用方」的场景。这时自定义一个 RejectedExecutionHandler,做两件事:计数上报(补上 JDK 不记的拒绝次数)和兜底落盘(拒绝的任务存起来,后台补偿重放)。
public class PersistingRejectedHandler implements RejectedExecutionHandler { private final LongAdder rejectedCount = new LongAdder(); private final TaskFallbackStore store; // 落库 / 发 MQ 都行 public PersistingRejectedHandler(TaskFallbackStore store) { this.store = store; } @Override public void rejectedExecution(Runnable r, ThreadPoolExecutor e) { rejectedCount.increment(); log.warn("任务被拒绝: task={}, active={}, queueSize={}", r, e.getActiveCount(), e.getQueue().size()); store.save(r); // 持久化,后台线程定时捞出重放 } // 暴露给监控系统(Micrometer / 自采) public long rejectedSum() { return rejectedCount.sum(); } }
三个使用细节:
- 计数用
LongAdder而不是AtomicLong——拒绝本身发生在高并发热点路径上,LongAdder 的分片累加竞争更小(呼应「CAS 与原子操作」篇)。 store.save(r)里如果任务不是可序列化的(比如捕获了栈上变量的 lambda),要改成把业务参数(orderId 之类)存下来,重放时重建任务。- 重放线程要限流,别把积压任务一次性灌回同一个池子,否则拒绝→落盘→重放→再拒绝,死循环。
可丢(日志/埋点)→ DiscardPolicy;要反压 → CallerRunsPolicy;只关心最新值 → DiscardOldestPolicy;不能丢且不能阻塞 → 自定义落盘;默认没想清楚 → AbortPolicy 让问题暴露出来,别用静默策略把问题藏起来。
优雅关闭:shutdown / shutdownNow / awaitTermination 三件套
发布、缩容、停服时线程池怎么关?直接让 JVM 退出会丢任务;关得太慢又卡住发布流程。标准做法是「两阶段关闭」,对应第 6 站的状态机:
// 阶段一:stop accepting(RUNNING → SHUTDOWN) // 不再接受新任务;队列里 + 执行中的任务继续跑完 pool.shutdown(); // 阶段二:给存量任务一个处理窗口 if (!pool.awaitTermination(30, TimeUnit.SECONDS)) { // 30s 还没跑完 → 中断运行中线程 + 取出未执行任务(→ STOP) List<Runnable> dropped = pool.shutdownNow(); if (!pool.awaitTermination(10, TimeUnit.SECONDS)) { log.error("线程池强制关闭失败,仍有线程存活"); } // shutdownNow() 返回的队列任务要自己兜底(落盘补偿) dropped.forEach(fallbackStore::save); } // 全部线程退出后 → TIDYING → TERMINATED,terminated() 钩子执行
| 方法 | 行为 | 新任务 | 存量任务 |
|---|---|---|---|
shutdown() | 状态 → SHUTDOWN | 拒绝(走拒绝策略) | 继续执行完毕 |
shutdownNow() | 状态 → STOP | 拒绝 | 尝试 interrupt 运行中线程,返回队列中未执行的任务 |
awaitTermination(t) | 阻塞等待池走到 TERMINATED | 超时返回 false,线程可继续 shutdownNow() | |
把两阶段关闭放进 @PreDestroy(或实现 DisposableBean),容器销毁时自动触发。注意 awaitTermination 的总时长要小于发布系统的优雅停机超时(通常 30~60s),否则发布流程会等超时后强杀进程,你精心写的兜底就白写了。
为什么 interrupt 不一定停得下来?
shutdownNow() 只是给运行中线程发 Thread.interrupt(),线程停不停取决于任务代码是否响应中断标志(Thread.sleep、BlockingQueue.take 这类会抛 InterruptedException 的操作会响应;纯 CPU 死循环不检查标志位就不响应)。所以任务代码里该检查 Thread.interrupted() 的长循环要检查——优雅关闭的上限,由你任务代码的中断响应能力决定。
生产级线程池模板:可直接复制
把前面 14 站的结论合成一个模板,以后新建线程池直接抄:
public class ThreadPoolTemplate { private static final int CPU_COUNT = Runtime.getRuntime().availableProcessors(); public static ThreadPoolExecutor create(String name, boolean ioBound) { int core = ioBound ? (int) (CPU_COUNT * (1 + 4.0)) // IO 密集,假设 W/C = 4 : CPU_COUNT + 1; // CPU 密集 int max = core * 2; int queueCap = ioBound ? 2000 : 500; ThreadPoolExecutor pool = new ThreadPoolExecutor( core, max, 60L, TimeUnit.SECONDS, new ArrayBlockingQueue<>(queueCap), new ThreadFactoryBuilder() .setNameFormat(name + "-%d") .setUncaughtExceptionHandler((t, e) -> log.error("Thread {} threw exception", t.getName(), e)) .build(), new ThreadPoolExecutor.CallerRunsPolicy() ); pool.prestartAllCoreThreads(); // 预热,避免首请求付建线程成本 return pool; } public static void shutdownGracefully(ThreadPoolExecutor pool) { pool.shutdown(); if (!pool.awaitTermination(30, TimeUnit.SECONDS)) { pool.shutdownNow(); } } }
为什么 ArrayBlockingQueue(有界,OOM 兜底)?为什么 CallerRunsPolicy(反压而不是静默丢)?为什么自定义 ThreadFactory(命名 + 异常兜底)?为什么 prestartAllCoreThreads(首包延迟)?每个选择都有上一站的推导,抄的时候记得把参数按自己业务的 W/C 和队列容量重新算一遍。
这一篇你掌握了什么
核心知识点回顾
- 7 大参数:corePoolSize / maximumPoolSize / keepAliveTime / unit / workQueue / threadFactory / rejectedExecutionHandler——餐厅模型一次讲清。
- execute() 决策树:核心线程 → 队列 → 非核心线程 → 拒绝;入队后 recheck 防 shutdown 竞态。
- ctl 位运算:高 3 位 5 状态(RUNNING/SHUTDOWN/STOP/TIDYING/TERMINATED)+ 低 29 位线程数,一个 AtomicInteger 原子地管理两件事。
- 队列选型:生产首选有界 ArrayBlockingQueue;无界队列是 OOM 温床;SynchronousQueue 纯并发。
- 4 种拒绝策略 + 自定义落盘策略(LongAdder 计数,补上 JDK 不记的拒绝次数)。
- 异常处理:execute 抛给 UncaughtExceptionHandler,submit 吞进 Future——两个兜底都要有。
- Executors 工厂禁用:无界队列 OOM / 无限线程爆炸,阿里规约【强制】。
- 调优:CPU 密集 N+1,IO 密集 N×(1+W/C);公式给起点,压测给答案。
- 优雅关闭:shutdown → awaitTermination → shutdownNow 两阶段,挂 @PreDestroy。
一句话总结
线程池的本质是一个「有界缓冲 + 弹性线程 + 明确拒绝」的流量整形器:参数决定它的形状,监控告诉你它现在什么状态,拒绝策略和优雅关闭决定它崩溃时体面不体面。
"线程池有 7 大参数,execute() 流程是先看核心线程、再看队列、再看最大线程、最后拒绝,状态和线程数打包在 ctl 这个 AtomicInteger 里——高 3 位 5 个生命周期状态,低 29 位线程数,CAS 保证原子性。生产环境禁止用 Executors 工厂方法:newFixedThreadPool 和 newSingleThreadExecutor 是无界队列会 OOM,newCachedThreadPool 线程数无上限会线程爆炸。推荐直接 new ThreadPoolExecutor,有界队列 + CallerRunsPolicy 做反压,IO 密集型按 N×(1+W/C) 定核心线程数,异常处理上 execute 靠 UncaughtExceptionHandler、submit 必须消费 Future 结果,关闭用 shutdown + awaitTermination 两阶段优雅停机,监控上盯 activeCount、队列水位和自定义 handler 统计的拒绝次数。"
Comments · 评论