这是《JVM 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 AQS 是 J.U.C 中锁和同步器的基础框架,线程池是 Java 应用并发执行的常用形态。理解它们,可以诊断 BLOCKED、任务堆积、拒绝执行、线程泄漏和 CPU 饱和。
12.1 AQS 结构
AbstractQueuedSynchronizer 核心思想:
state
+ FIFO wait queue
+ CAS
+ park / unpark
Thread A
CAS state 成功
-> 获取资源
Thread B
CAS state 失败
-> 入队
-> park
-> A 释放
-> unpark B
常见实现:
| 类 | 用途 |
|---|---|
| ReentrantLock | 可重入互斥锁 |
| Semaphore | 许可证 |
| CountDownLatch | 一次倒计数 |
| CyclicBarrier | 循环屏障 |
| ReentrantReadWriteLock | 读写锁 |
| ThreadPoolExecutor.Worker | 工作线程 |
12.2 ReentrantLock
ReentrantLock lock = new ReentrantLock();
try {
lock.lock();
// 临界区
} finally {
lock.unlock();
}
tryLock:
if (!lock.tryLock(100, TimeUnit.MILLISECONDS)) {
throw new TimeoutException("lock timeout");
}
公平锁:
ReentrantLock fair = new ReentrantLock(true);
公平锁降低饥饿概率但可能降低吞吐。多数业务使用非公平锁。
12.3 读写锁
ReentrantReadWriteLock lock = new ReentrantReadWriteLock();
lock.readLock().lock();
try {
// 读
} finally {
lock.readLock().unlock();
}
lock.writeLock().lock();
try {
// 写
} finally {
lock.writeLock().unlock();
}
适用:读多写少、读操作耗时、可接受一定协调开销。
StampedLock:
StampedLock lock = new StampedLock();
long stamp = lock.tryOptimisticRead();
int value = data;
if (!lock.validate(stamp)) {
stamp = lock.readLock();
try {
value = data;
} finally {
lock.unlockRead(stamp);
}
}
StampedLock 不可重入,写法复杂,必须有明确压测收益。
12.4 线程池参数
ThreadPoolExecutor executor = new ThreadPoolExecutor(
8,
16,
60L,
TimeUnit.SECONDS,
new ArrayBlockingQueue<>(1000),
new ThreadFactoryBuilder().setNameFormat("order-%d").build(),
new ThreadPoolExecutor.CallerRunsPolicy()
);
参数:
| 参数 | 含义 |
|---|---|
| corePoolSize | 核心线程 |
| maximumPoolSize | 最大线程 |
| keepAliveTime | 非核心线程空闲时间 |
| workQueue | 任务队列 |
| threadFactory | 线程创建 |
| rejectedExecutionHandler | 拒绝策略 |
12.5 提交流程
提交任务
-> 当前线程数 < core
-> 创建核心线程
-> 否则入队
-> 队列满且线程数 < max
-> 创建非核心线程
-> 队列满且线程数 = max
-> 拒绝
这意味着使用无界队列时,maximumPoolSize 不会发挥作用。
线程数估算:
CPU 密集:约等于 CPU 核数
IO 密集:核数 * (1 + 等待时间 / 计算时间)
公式只是起点,最终以压测和延迟目标为准。
12.6 拒绝策略
| 策略 | 行为 |
|---|---|
| AbortPolicy | 抛 RejectedExecutionException |
| CallerRunsPolicy | 提交线程执行 |
| DiscardPolicy | 静默丢弃 |
| DiscardOldestPolicy | 丢最老任务 |
生产建议:
- 显式设置有界队列;
- 拒绝策略记录日志和指标;
- 关键业务可降级或落盘;
- CallerRuns 可能拖慢上游;
- 静默丢弃不可接受。
自定义拒绝:
RejectedExecutionHandler handler = (r, executor1) -> {
log.error("task rejected, queue={}, pool={}",
executor1.getQueue().size(), executor1.getPoolSize());
throw new RejectedExecutionException("executor saturated");
};
12.7 监控线程池
executor.getActiveCount();
executor.getPoolSize();
executor.getQueue().size();
executor.getCompletedTaskCount();
executor.getTaskCount();
建议暴露指标:
active_count
pool_size
largest_pool_size
queue_size
completed_task_count
rejected_count
task_latency
异常场景:
| 现象 | 方向 |
|---|---|
| queue 持续增长 | 消费不足 |
| rejected 增长 | 容量不足或突发流量 |
| active == max 且 CPU 低 | 任务阻塞在外部 IO |
| 线程数持续增长 | 线程池泄漏 |
| CPU 高 | 任务计算密集或线程过多 |
12.8 优雅关闭
executor.shutdown();
if (!executor.awaitTermination(30, TimeUnit.SECONDS)) {
List<Runnable> pending = executor.shutdownNow();
log.warn("forced shutdown, pending={}", pending.size());
}
区别:
| 方法 | 行为 |
|---|---|
| shutdown | 不接新任务,存量继续 |
| shutdownNow | 尝试中断并返回等待任务 |
| awaitTermination | 等待结束 |
任务要正确响应中断:
try {
while (!Thread.currentThread().isInterrupted()) {
process();
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
12.9 常见错误
错误一:Executors 隐藏容量:
Executors.newFixedThreadPool(10);
Executors.newCachedThreadPool();
固定线程池通常是无界 LinkedBlockingQueue,缓存线程池最大线程数近似无界。
正确做法是显式 new ThreadPoolExecutor 并指定容量。
错误二:没有命名:
Executors.defaultThreadFactory()
排查线程 dump 时很难定位。应使用自定义 ThreadFactory。
错误三:提交任务无超时:
Future<String> future = executor.submit(task);
future.get();
应设置:
future.get(3, TimeUnit.SECONDS);
错误四:重复创建线程池:
每个请求创建一个线程池会导致线程和内存泄漏。线程池应生命周期明确、单例复用并统一治理。
12.10 CompletableFuture
CompletableFuture<Order> orderFuture =
CompletableFuture.supplyAsync(this::loadOrder, executor);
CompletableFuture<User> userFuture =
CompletableFuture.supplyAsync(this::loadUser, executor);
OrderView view = orderFuture
.thenCombine(userFuture, this::merge)
.orTimeout(2, TimeUnit.SECONDS)
.join();
注意:
- 显式传入业务线程池;
- 设置超时;
- 处理异常;
- 避免默认 ForkJoinPool 承担远程 IO;
- 线程上下文要传递;
- 记录 trace。
12.11 线程 dump 中的线程池
jstack <pid> > thread.txt
jcmd <pid> Thread.print -l > thread.txt
观察:
"order-3" #125 daemon prio=5 ...
java.lang.Thread.State: WAITING (parking)
at jdk.internal.misc.Unsafe.park
at java.util.concurrent.locks.LockSupport.park
判断:
- 线程数量;
- 线程名称;
- 等待位置;
- 锁持有者;
- 是否重复创建;
- 是否阻塞外部调用。
本章小结
AQS 通过 state、队列、CAS 和 park 实现同步器。线程池要显式设置核心数、最大线程、有界队列、命名工厂、拒绝策略和监控。优雅停机依赖 shutdown、awaitTermination、任务可中断和上下游摘流。
思考题
- AQS 的核心组成是什么?
- 线程池提交任务的顺序是什么?
- 为什么不推荐无界队列?
- 四种拒绝策略分别适合什么场景?
- 如何为一个 IO 密集服务设计线程池?