JVMNotes

第 12 章:AQS 与线程池

zjc 于 2026-01-12 发布

这是《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 丢最老任务

生产建议:

  1. 显式设置有界队列;
  2. 拒绝策略记录日志和指标;
  3. 关键业务可降级或落盘;
  4. CallerRuns 可能拖慢上游;
  5. 静默丢弃不可接受。

自定义拒绝:

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();

注意:

  1. 显式传入业务线程池;
  2. 设置超时;
  3. 处理异常;
  4. 避免默认 ForkJoinPool 承担远程 IO;
  5. 线程上下文要传递;
  6. 记录 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

判断:

  1. 线程数量;
  2. 线程名称;
  3. 等待位置;
  4. 锁持有者;
  5. 是否重复创建;
  6. 是否阻塞外部调用。

本章小结

AQS 通过 state、队列、CAS 和 park 实现同步器。线程池要显式设置核心数、最大线程、有界队列、命名工厂、拒绝策略和监控。优雅停机依赖 shutdown、awaitTermination、任务可中断和上下游摘流。

思考题

  1. AQS 的核心组成是什么?
  2. 线程池提交任务的顺序是什么?
  3. 为什么不推荐无界队列?
  4. 四种拒绝策略分别适合什么场景?
  5. 如何为一个 IO 密集服务设计线程池?