线程池详解
飞翔科技会议室,深夜十一点。
大翔(CTO,盯着监控大屏):"系统吞吐量又跌了30%!线上每来一个请求就 new 一个 Thread,十分钟飙到两万个线程,CPU 全耗在上下文切换上了。"
白歌(首席架构师,推了推眼镜):"典型的线程泛滥。JDK 8 的 ThreadPoolExecutor 就是为解决这个问题设计的——线程复用、任务队列、拒绝策略,七个参数一把梭。"
小崔(初级开发,挠头):"七个参数...我只记得 newFixedThreadPool 和 newCachedThreadPool,其他都是啥啊?"
朱璐(测试组长):"上次压测你们用的 newCachedThreadPool,TPS 冲到 5000 就 OOM 了,SynchronousQueue 无界创建线程,服务器直接冒烟。"
Frank(后端专家):"Let me explain. Thread pool is not just about creating threads — it's about managing resources. When you understand all seven parameters, you can tune the pool to exactly match your workload."
大翔:"好,今晚就把线程池从原理到实战彻底讲透!"
一、定义表
| 概念 | 定义 | JDK 8 关键类/接口 |
|---|---|---|
| 线程池 | 预先创建并维护一组工作线程,复用线程执行提交的任务,避免频繁创建/销毁线程的开销 | ThreadPoolExecutor |
| 核心线程数 | 线程池中始终保持存活的线程数,即使处于空闲状态也不会被回收 | corePoolSize |
| 最大线程数 | 线程池中允许的最大线程数,当队列满时才会创建超出核心数的线程 | maximumPoolSize |
| 空闲超时 | 非核心线程空闲超过此时间后会被回收 | keepAliveTime |
| 工作队列 | 用于缓存待执行任务的阻塞队列 | BlockingQueue<Runnable> |
| 线程工厂 | 自定义线程创建逻辑(命名、优先级、守护状态) | ThreadFactory |
| 拒绝策略 | 当线程池和队列都满时,对新提交任务的处理方式 | RejectedExecutionHandler |
| Future/Callable | 可返回结果、可感知异常的任务抽象 | Future<T>, Callable<T> |
| CompletionService | 解耦任务提交与结果获取,按完成顺序获取结果 | ExecutorCompletionService<T> |
二、ThreadPoolExecutor 七个参数深度解析
2.1 构造函数签名
public ThreadPoolExecutor(
int corePoolSize, // 核心线程数
int maximumPoolSize, // 最大线程数
long keepAliveTime, // 空闲线程存活时间
TimeUnit unit, // 时间单位
BlockingQueue<Runnable> workQueue, // 工作队列
ThreadFactory threadFactory, // 线程工厂
RejectedExecutionHandler handler // 拒绝策略
)
2.2 七个参数协同原理
2.3 线程池状态机
三、Executors 工厂方法
| 工厂方法 | 核心线程 | 最大线程 | 队列类型 | 适用场景 |
|---|---|---|---|---|
newFixedThreadPool(n) | n | n | LinkedBlockingQueue(无界) | 负载稳定、并发数可控 |
newCachedThreadPool() | 0 | Integer.MAX_VALUE | SynchronousQueue(不存储) | 短期异步任务、突发流量 |
newSingleThreadExecutor() | 1 | 1 | LinkedBlockingQueue(无界) | 任务串行执行、保证顺序 |
newScheduledThreadPool(n) | n | Integer.MAX_VALUE | DelayedWorkQueue | 定时/周期性任务 |
白歌提醒:
newCachedThreadPool的最大线程数是Integer.MAX_VALUE,高并发下可能创建几十万个线程导致 OOM。生产环境强烈建议手动构造ThreadPoolExecutor。
四、四种拒绝策略
| 策略 | 行为 | 适用场景 |
|---|---|---|
AbortPolicy(默认) | 抛出 RejectedExecutionException | 需要立即感知任务被拒绝 |
CallerRunsPolicy | 由提交任务的线程(调用者)同步执行该任务 | 需要降级处理,让调用者分担压力 |
DiscardPolicy | 静默丢弃被拒绝的任务,不抛异常 | 允许丢失部分任务(如日志采集) |
DiscardOldestPolicy | 丢弃队列中最老的任务,然后重试提交当前任务 | 优先处理最新数据 |
五、代码示例
示例一:飞翔科技批量处理员工数据——使用自定义 ThreadPoolExecutor
场景:飞翔科技人力资源系统需要在年终批量计算 100 名员工的绩效奖金,每名员工的计算涉及数据库查询和复杂公式运算,为控制并发使用自定义线程池。
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;
/**
* 飞翔科技 —— 员工年终绩效奖金批量计算
*/
public class EmployeeBonusCalculator {
// 自定义线程工厂:给线程起有意义的名字
static class FeixiangThreadFactory implements ThreadFactory {
private final AtomicInteger counter = new AtomicInteger(1);
private final String prefix;
public FeixiangThreadFactory(String prefix) {
this.prefix = prefix;
}
@Override
public Thread newThread(Runnable r) {
Thread t = new Thread(r, prefix + "-" + counter.getAndIncrement());
t.setDaemon(false); // 非守护线程,保证任务执行完
t.setPriority(Thread.NORM_PRIORITY);
return t;
}
}
// 模拟员工绩效计算任务
static class BonusTask implements Callable<String> {
private final int employeeId;
private final String employeeName;
public BonusTask(int employeeId, String employeeName) {
this.employeeId = employeeId;
this.employeeName = employeeName;
}
@Override
public String call() throws Exception {
// 模拟耗时计算
long computeTime = 200 + (long)(Math.random() * 800);
Thread.sleep(computeTime);
double bonus = 5000 + Math.random() * 15000;
return String.format("员工[%d] %s: 年终奖金 ¥%.2f (计算耗时 %dms)",
employeeId, employeeName, bonus, computeTime);
}
}
public static void main(String[] args) {
// 手动构造 ThreadPoolExecutor —— 七个参数全掌控
ThreadPoolExecutor executor = new ThreadPoolExecutor(
4, // corePoolSize: 核心4个线程常驻
8, // maximumPoolSize: 最多8个线程
60L, // keepAliveTime: 空闲60秒
TimeUnit.SECONDS, // 时间单位
new LinkedBlockingQueue<>(20), // 有界队列,容量20
new FeixiangThreadFactory("bonus-worker"), // 自定义线程工厂
new ThreadPoolExecutor.CallerRunsPolicy() // 拒绝策略:调用者运行
);
// 允许核心线程超时回收(JDK 8 新方法)
executor.allowCoreThreadTimeOut(false);
// 提交 100 个计算任务
CompletionService<String> completionService =
new ExecutorCompletionService<>(executor);
String[] names = {"张伟", "李娜", "王磊", "陈静", "赵强", "孙莉", "周杰", "吴敏",
"郑浩", "钱雪", "冯涛", "蒋琳", "韩冰", "杨帆", "朱峰", "秦月",
"尤佳", "许明", "何琳", "吕华"};
System.out.println("===== 飞翔科技年终绩效计算开始 =====");
System.out.println("线程池状态 - 核心: " + executor.getCorePoolSize()
+ ", 最大: " + executor.getMaximumPoolSize()
+ ", 队列容量: 20");
System.out.println("共需计算员工数: 100");
for (int i = 0; i < 100; i++) {
final int employeeId = 1001 + i;
final String name = names[i % names.length] + (i / names.length + 1);
completionService.submit(new BonusTask(employeeId, name));
}
// 按完成顺序获取结果
System.out.println("\n===== 计算结果(按完成先后顺序输出)=====");
for (int i = 0; i < 100; i++) {
try {
Future<String> future = completionService.take();
System.out.println("[" + (i + 1) + "] " + future.get());
} catch (InterruptedException | ExecutionException e) {
System.err.println("计算失败: " + e.getMessage());
}
}
System.out.println("\n===== 线程池统计 =====");
System.out.println("已完成任务数: " + executor.getCompletedTaskCount());
System.out.println("历史最大线程数: " + executor.getLargestPoolSize());
System.out.println("当前活跃线程数: " + executor.getActiveCount());
// 优雅关闭
executor.shutdown();
try {
if (!executor.awaitTermination(30, TimeUnit.SECONDS)) {
executor.shutdownNow();
}
} catch (InterruptedException e) {
executor.shutdownNow();
}
System.out.println("线程池已关闭,程序结束。");
}
}
控制台输出(部分):
===== 飞翔科技年终绩效计算开始 =====
线程池状态 - 核心: 4, 最大: 8, 队列容量: 20
共需计算员工数: 100
===== 计算结果(按完成先后顺序输出)=====
[1] 员工[1051] 张伟6: 年终奖金 ¥18234.56 (计算耗时 312ms)
[2] 员工[1003] 王磊1: 年终奖金 ¥7654.12 (计算耗时 245ms)
[3] 员工[1017] 孙莉4: 年终奖金 ¥15678.90 (计算耗时 567ms)
...
[100] 员工[1092] 吕华5: 年终奖金 ¥9876.54 (计算耗时 723ms)
===== 线程池统计 =====
已完成任务数: 100
历史最大线程数: 8
当前活跃线程数: 0
线程池已关闭,程序结束。
示例二:定时任务调度——飞翔科技每日系统健康检查
场景:运维组需要两个定时任务:① 每 5 秒采集一次服务器健康指标;② 每天凌晨 2 点执行数据库清理。使用 ScheduledThreadPoolExecutor 实现。
import java.util.concurrent.*;
import java.time.LocalDateTime;
import java.time.format.DateTimeFormatter;
/**
* 飞翔科技 —— 定时系统健康检查与维护
*/
public class SystemHealthScheduler {
private static final DateTimeFormatter fmt =
DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss");
// 模拟健康指标采集
static class HealthCollectTask implements Runnable {
private final int serverId;
public HealthCollectTask(int serverId) {
this.serverId = serverId;
}
@Override
public void run() {
double cpuUsage = 20 + Math.random() * 60;
double memoryUsage = 30 + Math.random() * 50;
String status = cpuUsage > 80 ? "⚠ 警告" : "✓ 正常";
System.out.printf("[%s] 服务器 #%d | CPU: %.1f%% | 内存: %.1f%% | %s%n",
LocalDateTime.now().format(fmt), serverId,
cpuUsage, memoryUsage, status);
}
}
// 模拟数据库清理任务
static class DatabaseCleanupTask implements Runnable {
@Override
public void run() {
System.out.println("\n========================================");
System.out.printf("[%s] 开始执行数据库过期日志清理...%n",
LocalDateTime.now().format(fmt));
System.out.println("清理 30 天前的操作日志...");
System.out.println("清理 90 天前的审计记录...");
System.out.println("重建索引...");
System.out.println("数据库清理完成!");
System.out.println("========================================\n");
}
}
public static void main(String[] args) throws InterruptedException {
ScheduledThreadPoolExecutor scheduler = new ScheduledThreadPoolExecutor(
3,
new ThreadFactory() {
private int count = 0;
@Override
public Thread newThread(Runnable r) {
return new Thread(r, "scheduler-" + (++count));
}
},
new ThreadPoolExecutor.DiscardPolicy() // 任务过多时丢弃
);
System.out.println("===== 飞翔科技系统健康监控启动 =====");
System.out.println("当前时间: " + LocalDateTime.now().format(fmt));
// 任务1:初始延迟1秒,每5秒执行一次(采集3台服务器指标)
scheduler.scheduleAtFixedRate(
new HealthCollectTask(1), 1, 5, TimeUnit.SECONDS);
scheduler.scheduleAtFixedRate(
new HealthCollectTask(2), 2, 5, TimeUnit.SECONDS);
// 任务2:延迟执行(模拟每天凌晨2点清理——这里用10秒后执行做演示)
scheduler.schedule(new DatabaseCleanupTask(), 10, TimeUnit.SECONDS);
// 运行20秒后关闭
Thread.sleep(20000);
System.out.println("\n调度器关闭中...");
scheduler.shutdown();
if (!scheduler.awaitTermination(5, TimeUnit.SECONDS)) {
scheduler.shutdownNow();
}
System.out.println("调度器已关闭。");
}
}
控制台输出(部分):
===== 飞翔科技系统健康监控启动 =====
当前时间: 2025-01-15 14:30:00
[2025-01-15 14:30:01] 服务器 #1 | CPU: 45.2% | 内存: 62.8% | ✓ 正常
[2025-01-15 14:30:02] 服务器 #2 | CPU: 33.7% | 内存: 51.1% | ✓ 正常
[2025-01-15 14:30:06] 服务器 #1 | CPU: 72.1% | 内存: 78.4% | ✓ 正常
[2025-01-15 14:30:07] 服务器 #2 | CPU: 41.6% | 内存: 55.3% | ✓ 正常
========================================
[2025-01-15 14:30:10] 开始执行数据库过期日志清理...
清理 30 天前的操作日志...
清理 90 天前的审计记录...
重建索引...
数据库清理完成!
========================================
[2025-01-15 14:30:11] 服务器 #1 | CPU: 88.3% | 内存: 67.2% | ⚠ 警告
...
调度器关闭中...
调度器已关闭。
六、execute() vs submit() 对比
| 特性 | execute(Runnable) | submit(Callable<T>) / submit(Runnable) |
|---|---|---|
| 返回值 | 无(void) | 返回 Future<T> |
| 异常处理 | 异常抛出到 UncaughtExceptionHandler | 异常封装在 Future.get() 中 |
| 适用场景 | 不需要返回结果,不关心异常 | 需要获取结果或感知异常 |
// execute 方式
executor.execute(() -> {
int result = 1 / 0; // 异常会被 UncaughtExceptionHandler 捕获
});
// submit 方式
Future<Integer> future = executor.submit(() -> {
return 1 / 0; // 异常在 future.get() 时抛出 ExecutionException
});
try {
future.get();
} catch (ExecutionException e) {
System.out.println("捕获到任务异常: " + e.getCause());
}
七、易错场景
7.1 错误配置:核心线程数大于最大线程数
// ❌ 错误:corePoolSize > maximumPoolSize 会抛出 IllegalArgumentException
ThreadPoolExecutor pool = new ThreadPoolExecutor(
10, // corePoolSize = 10
5, // maximumPoolSize = 5 ← 小于 corePoolSize!
60, TimeUnit.SECONDS,
new LinkedBlockingQueue<>(10),
Executors.defaultThreadFactory(),
new ThreadPoolExecutor.AbortPolicy()
);
// 运行结果:java.lang.IllegalArgumentException
正确做法:corePoolSize <= maximumPoolSize,通常 maximumPoolSize 至少为核心数的 1.5~2 倍。
7.2 使用无界队列导致 OOM
// ❌ 危险:LinkedBlockingQueue 无参构造 = 无界队列(容量 Integer.MAX_VALUE)
ThreadPoolExecutor pool = new ThreadPoolExecutor(
5, 10, 60, TimeUnit.SECONDS,
new LinkedBlockingQueue<>(), // 无界队列!任务可以无限堆积
Executors.defaultThreadFactory(),
new ThreadPoolExecutor.AbortPolicy()
);
// 高并发下任务堆积到内存耗尽 → OutOfMemoryError
for (int i = 0; i < 10_000_000; i++) {
pool.execute(() -> {
try { Thread.sleep(10000); } catch (InterruptedException e) {}
});
}
正确做法:使用有界队列,如 new LinkedBlockingQueue<>(1000) 或 new ArrayBlockingQueue<>(500)。
7.3 忘记调用 shutdown() 导致 JVM 无法退出
// ❌ 错误:线程池未关闭,非守护线程存活导致 JVM 挂起
ThreadPoolExecutor pool = new ThreadPoolExecutor(
2, 4, 60, TimeUnit.SECONDS,
new LinkedBlockingQueue<>(10),
new ThreadFactory() {
public Thread newThread(Runnable r) {
Thread t = new Thread(r);
t.setDaemon(false); // 非守护线程
return t;
}
},
new ThreadPoolExecutor.AbortPolicy()
);
pool.execute(() -> System.out.println("任务执行完毕"));
// 忘记调用 pool.shutdown() → JVM 不会退出!
正确做法:使用 try-finally 或在 finally 块中调用 shutdown() + awaitTermination()。
7.4 shutdown() 与 shutdownNow() 混用导致任务丢失
// ❌ 错误:shutdownNow() 会返回未执行的任务列表,不处理就会丢失
List<Runnable> abandoned = pool.shutdownNow();
// abandoned 包含所有尚未开始执行的任务,直接忽略则丢失
System.out.println("丢弃了 " + abandoned.size() + " 个任务"); // 应该重新提交或持久化
正确做法:对 shutdownNow() 返回的未执行任务列表进行处理——记录日志、持久化到数据库或重新提交。
八、面试考点
Q1:ThreadPoolExecutor 的任务处理流程是怎样的?(高频)
答:当提交一个新任务时,ThreadPoolExecutor 按以下顺序处理:
- 如果当前运行线程数 <
corePoolSize,直接创建新线程执行任务(即使有空闲核心线程也会创建,直到达到 corePoolSize)。 - 如果当前运行线程数 >=
corePoolSize,尝试将任务放入workQueue。 - 如果
workQueue已满,且当前运行线程数 <maximumPoolSize,创建新线程执行任务。 - 如果
workQueue已满,且当前运行线程数 >=maximumPoolSize,执行RejectedExecutionHandler的拒绝策略。
JDK 8 还新增了 allowCoreThreadTimeOut(true) 可以让核心线程在空闲时也超时回收。
Q2:为什么不建议用 Executors 工厂方法创建线程池?(阿里规约/高频)
答:
newFixedThreadPool(n)和newSingleThreadExecutor()使用无界队列LinkedBlockingQueue,任务堆积会导致 OOM。newCachedThreadPool()的maximumPoolSize = Integer.MAX_VALUE,无限创建线程导致 OOM。newScheduledThreadPool(n)的maximumPoolSize也是Integer.MAX_VALUE。- 最佳实践:手动构造
ThreadPoolExecutor,使用有界队列,明确拒绝策略,自定义线程工厂(命名线程便于排查)。
Q3:CallerRunsPolicy 的原理和优缺点是什么?
答:当线程池和队列都满时,CallerRunsPolicy 让提交任务的调用线程(caller)来同步执行该任务。
- 原理:
rejectedExecution()方法中直接调用r.run(),阻塞调用者。 - 优点:天然的流量削峰——调用者被阻塞后无法继续提交新任务,形成背压(backpressure);任务不会丢失。
- 缺点:如果调用者是主线程或关键线程,阻塞会影响整体吞吐量;不适合对延迟敏感的场景。
Q4:核心线程数为 0 的线程池有什么特点?JDK 8 中 allowCoreThreadTimeOut 的作用?
答:
corePoolSize = 0的线程池(如newCachedThreadPool)在没有任务时不会保留任何线程,所有线程都是临时的,空闲超时后全部回收。- JDK 8 中
allowCoreThreadTimeOut(true)允许核心线程也在空闲超过keepAliveTime后被回收。这在需要弹性伸缩、资源敏感的场景下非常有用。 - 需要注意:设置
allowCoreThreadTimeOut(true)后,corePoolSize参数实际上不再有"永久保留"的语义。
九、JDK 8 线程池最佳实践总结
- 永远手动构造
ThreadPoolExecutor,禁止使用Executors工厂方法。 - 有界队列 + 合理拒绝策略:
ArrayBlockingQueue或指定容量的LinkedBlockingQueue。 - 自定义 ThreadFactory:给线程起业务名,方便 jstack 排查。
- 核心线程数公式:CPU 密集型 →
N+1;IO 密集型 →2N或N * (1 + WT/ST)(N = CPU 核数,WT = 等待时间,ST = 计算时间)。 - 优雅关闭:
shutdown()→awaitTermination()→ 超时则shutdownNow()。 - 监控:定期检查
getActiveCount()、getQueue().size()、getCompletedTaskCount()。
下一篇:并发工具类 —— CountDownLatch、CyclicBarrier、Semaphore 等 JDK 8 并发利器实战。