Lesson 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 个线程处理不过来,任务全部堆在队列里排队,每个请求都在等队列里的任务被执行,而队列里的任务在等线程来取——死等,超时,雪崩。
线程池是并发编程面试出现频率最高的主题,没有之一。原因有三:
- 覆盖面广——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(所有线程都意外退出了),就创建一个非核心线程来保证队列中的任务能被消费。
workQueue 选择:4 种队列的取舍
workQueue 的选择直接决定了线程池在流量高峰期的表现。这是很多候选人容易忽略的点——他们能背出 7 个参数名,却说不出不同队列的适用场景。
| 队列类型 | 容量 | 特点 | 风险 | 适用场景 |
|---|---|---|---|---|
LinkedBlockingQueue |
默认 Integer.MAX_VALUE(≈无界) | 链表实现,吞吐量高 | OOM ! | 必须指定容量后使用 |
ArrayBlockingQueue |
有界,必须指定 capacity | 数组实现,一把 ReentrantLock | 队列满时触发 max 线程或拒绝 | 生产环境首选 |
SynchronousQueue |
0(不存储任务) | 直接交付,生产者必须等消费者 | 无线程接收就立即走拒绝策略 | Executors.newCachedThreadPool() |
PriorityBlockingQueue |
无界 | 按优先级排序 | 优先级低的任务可能饥饿 | 任务有明确优先级的场景 |
// 推荐:有界队列 + 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,但实现更简单、内存占用更低,且强制有界——对于线程池场景,这点吞吐量差异远不如"有界"带来的安全性重要。
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 优雅得多。
"生产环境一般用 CallerRunsPolicy 做反压。如果任务不能丢且不能阻塞调用者,我会自定义 RejectedExecutionHandler,把被拒绝的任务持久化到数据库或 MQ,事后补偿。"
生产环境参数调优:公式 + 实战
面试中被问"线程池参数怎么设",回答"看情况"是不够的。你需要给出公式、给出数字、给出监控手段。
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() );
公式只是起点,真正上线后必须配合监控数据持续调优:
// 定时采集,上报到 Prometheus / Grafana ScheduledExecutorService monitor = ...; monitor.scheduleAtFixedRate(() -> { log.info("[pool-monitor] active={}, poolSize={}, queue={}, completed={}", pool.getActiveCount(), // 正在执行任务的线程数 pool.getPoolSize(), // 当前池中线程总数 pool.getQueue().size(), // 队列中待执行的任务数 pool.getCompletedTaskCount() // 已完成的任务总数 ); }, 0, 10, TimeUnit.SECONDS);
| 监控指标 | 告警阈值 | 说明 |
|---|---|---|
getActiveCount() / getPoolSize() | > 80% | 线程利用率过高,考虑扩容 |
getQueue().size() / capacity | > 70% | 队列堆积,有拒绝风险 |
getRejectedExecutionCount() | > 0 | 出现拒绝,必须立即处理 |
getCompletedTaskCount() 增量 | 突降 | 吞吐量下降,可能有死锁或下游故障 |
真实调优经历:公式算出 40,实际只用了 25
某服务按公式设了 corePoolSize=40,但上线后发现线程切换开销明显——CPU 利用率只有 60% 但 load average 很高。原因:下游 RPC 超时设得很短(200ms),大多数 IO 等待时间其实没有 40ms 那么长。调低到 25 后,吞吐量反而提升了 15%。结论:公式给起点,监控给方向,压测给答案。
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 后问题解决。
总结:参数速查表 + 优雅关闭
7 大参数速查
| 参数 | 一句话 | 生产建议 |
|---|---|---|
| corePoolSize | 常驻线程数 | CPU 密集: N+1 / IO 密集: N*(1+W/C) |
| maximumPoolSize | 线程上限 | core 的 1.5~2 倍,留弹性 |
| keepAliveTime | 非核心线程空闲存活时间 | 60s,高波动场景可设 30s |
| unit | 时间单位 | SECONDS |
| workQueue | 任务等待队列 | ArrayBlockingQueue + 合理容量 |
| threadFactory | 线程工厂 | 必须设置有意义的线程名 |
| rejectedHandler | 拒绝策略 | CallerRunsPolicy 或自定义持久化 |
// shutdown() —— 温和关闭 // 1. 不再接受新任务 // 2. 已提交的任务(队列中 + 执行中)继续执行完毕 // 3. 所有任务完成后,线程自然退出 pool.shutdown(); if (!pool.awaitTermination(30, TimeUnit.SECONDS)) { pool.shutdownNow(); // 30s 还没关完 → 强制中断 } // shutdownNow() —— 暴力关闭 // 1. 不再接受新任务 // 2. 尝试中断所有正在执行的线程(Thread.interrupt()) // 3. 返回队列中尚未执行的任务列表 List<Runnable> pending = pool.shutdownNow();
生产环境推荐使用"两阶段关闭":先调 shutdown() 给线程池一个处理剩余任务的机会,然后通过 awaitTermination() 等待一段时间。如果超时仍未关闭,再调 shutdownNow() 强制中断。在 Spring 环境中,可以通过 @PreDestroy 或 DisposableBean 实现这一逻辑。
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; return 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() ); } }
"线程池有 7 大参数,execute() 流程是先看核心线程、再看队列、再看最大线程、最后拒绝。生产环境禁止使用 Executors 工厂方法,因为 newFixedThreadPool 和 newSingleThreadExecutor 的无界队列会导致 OOM,newCachedThreadPool 的无限线程会导致线程爆炸。推荐直接 new ThreadPoolExecutor,配合有界队列和 CallerRunsPolicy,根据 CPU/IO 密集型选择合理的核心线程数,并通过 getActiveCount、getQueue().size() 等指标做持续监控。"