并发编程总结
开篇:并发编程的全景图
前面几篇我们分别聊了线程基础、synchronized、volatile、AQS、原子类、ThreadLocal。现在是时候把这些零散的知识串成一张网了。
本文是并发编程系列的"大总结",覆盖线程池、Fork/Join、并发容器、阻塞队列、异步编程等实战中最常用的工具。每个工具只讲核心原理和使用姿势,不再逐行贴源码,需要深挖源码的地方会指向对应的专题文章。
一、线程池深入
1.1 为什么需要线程池?
线程的创建和销毁都是昂贵的操作(涉及 OS 线程调度、内存分配)。如果每个任务都 new 一个线程,高并发下系统很快就会 OOM。线程池的核心思想是复用:提前创建好一批线程放在池子里,任务来了就从池子里取线程执行,执行完线程归还池子,而不是销毁。
1.2 ThreadPoolExecutor 的 7 个参数
这是线程池最重要的知识点,面试必问。
public ThreadPoolExecutor(
int corePoolSize, // 核心线程数(常驻员工)
int maximumPoolSize, // 最大线程数(常驻 + 临时员工)
long keepAliveTime, // 非核心线程空闲存活时间
TimeUnit unit, // 存活时间单位
BlockingQueue<Runnable> workQueue, // 任务队列
ThreadFactory threadFactory, // 线程工厂
RejectedExecutionHandler handler // 拒绝策略
)1.3 任务提交流程
当你调用 execute(task) 时,线程池内部的决策流程如下:
提交任务
│
▼
当前线程数 < corePoolSize?──── 是 ──→ 创建核心线程执行
│ 否
▼
任务队列未满?──── 是 ──→ 任务入队等待
│ 否
▼
当前线程数 < maximumPoolSize?── 是 ──→ 创建非核心线程执行
│ 否
▼
执行拒绝策略注意顺序:先核心线程,再队列,最后非核心线程。这个顺序说明线程池更倾向于让任务排队而不是创建新线程。
1.4 四种拒绝策略
| 策略 | 行为 |
|---|---|
| AbortPolicy(默认) | 抛出 RejectedExecutionException |
| DiscardPolicy | 默默丢弃,不抛异常 |
| DiscardOldestPolicy | 丢弃队列中最老的任务,提交新任务 |
| CallerRunsPolicy | 让提交任务的线程自己执行 |
生产环境最常用的是 AbortPolicy(让调用方感知到问题)和 CallerRunsPolicy(降速不丢任务)。
1.5 线程池状态
线程池用 ctl(一个 AtomicInteger)的高 3 位表示状态:
| 状态 | 含义 |
|---|---|
| RUNNING | 正常接收和处理任务 |
| SHUTDOWN | 不接新任务,处理完已有任务 |
| STOP | 不接新任务,中断正在处理的任务 |
| TIDYING | 所有任务已终止,正在做清理 |
| TERMINATED | 完全终止 |
状态转换:RUNNING → SHUTDOWN → TIDYING → TERMINATED,或 RUNNING → STOP → TIDYING → TERMINATED。
1.6 Executors 提供的线程池
| 工厂方法 | 核心线程 | 最大线程 | 队列 | 特点 |
|---|---|---|---|---|
| newFixedThreadPool(n) | n | n | LinkedBlockingQueue(无界) | 固定线程数 |
| newSingleThreadExecutor | 1 | 1 | LinkedBlockingQueue(无界) | 单线程顺序执行 |
| newCachedThreadPool | 0 | Integer.MAX_VALUE | SynchronousQueue | 来一个任务就开一个线程 |
| newScheduledThreadPool(n) | n | Integer.MAX_VALUE | DelayedWorkQueue | 定时/周期执行 |
为什么阿里规范禁止用 Executors? 因为 newFixedThreadPool 和 newSingleThreadExecutor 的队列是无界的(Integer.MAX_VALUE),可能堆积大量任务导致 OOM;newCachedThreadPool 的最大线程数是无界的,可能创建大量线程导致 OOM。
正确做法是自己用 ThreadPoolExecutor 构造,给队列指定容量:
ExecutorService pool = new ThreadPoolExecutor(
10, 20, 60, TimeUnit.SECONDS,
new ArrayBlockingQueue<>(1024),
new ThreadPoolExecutor.AbortPolicy()
);1.7 线程数设定
没有万能公式,但有经验参考:
- CPU 密集型:线程数 = CPU 核数 + 1
- IO 密集型:线程数 = CPU 核数 x 2 + 1
更精确的公式:线程数 = CPU核数 x 目标CPU利用率 x (1 + 等待时间/计算时间)
但实际生产中,最靠谱的方法是:先按公式给个初始值,然后通过压测不断调整,找到最佳平衡点。
1.8 动态线程池
运行时通过 setCorePoolSize() 和 setMaximumPoolSize() 可以动态调整线程数,调用后立即生效。结合监控指标(队列饱和度、线程活跃度、CPU 利用率等),可以实现根据负载自动调整的动态线程池。开源方案有 DynamicTp 等。
1.9 ScheduledThreadPoolExecutor
用于执行定时任务的线程池,核心原理是 DelayedWorkQueue(基于堆的延迟队列)。
ScheduledThreadPoolExecutor pool = new ScheduledThreadPoolExecutor(5);
// 延迟 2 秒执行
pool.schedule(() -> doTask(), 2, TimeUnit.SECONDS);
// 固定频率:每 3 秒执行一次(从任务开始计算)
pool.scheduleAtFixedRate(() -> doTask(), 2, 3, TimeUnit.SECONDS);
// 固定延迟:上一次任务结束后等 3 秒再执行
pool.scheduleWithFixedDelay(() -> doTask(), 2, 3, TimeUnit.SECONDS);scheduleAtFixedRate 和 scheduleWithFixedDelay 的区别:前者从任务开始计时,后者从任务结束计时。如果任务执行时间超过周期间隔,scheduleAtFixedRate 不会并行执行,而是任务结束后立即开始下一次。
二、Fork/Join 框架
2.1 核心思想:分而治之 + 工作窃取
Fork/Join 框架是 JDK 7 引入的并行计算框架,核心思想是分而治之:把大任务拆成小任务(Fork),各线程并行执行,最后合并结果(Join)。
它和 ThreadPoolExecutor 最大的区别是工作窃取算法:每个线程有自己的任务双端队列,空闲线程会从其他线程的队列尾部偷任务来做。这样就不会出现有的线程忙死、有的线程闲死的情况。
Thread-1 队列: [A1, A2, A3] ← 自己从头部取
Thread-2 队列: [B1]
Thread-3 队列: [] ──────→ 从 Thread-1 尾部偷 A32.2 使用示例
class SumTask extends RecursiveTask<Long> {
private int[] arr;
private int lo, hi;
SumTask(int[] arr, int lo, int hi) {
this.arr = arr; this.lo = lo; this.hi = hi;
}
@Override
protected Long compute() {
if (hi - lo <= 1000) {
// 小任务直接计算
long sum = 0;
for (int i = lo; i < hi; i++) sum += arr[i];
return sum;
}
int mid = (lo + hi) >>> 1;
SumTask left = new SumTask(arr, lo, mid);
SumTask right = new SumTask(arr, mid, hi);
left.fork(); // 提交左半到队列
long rResult = right.compute(); // 右半自己算
long lResult = left.join(); // 等左半结果
return lResult + rResult;
}
}
ForkJoinPool pool = new ForkJoinPool();
long result = pool.invoke(new SumTask(bigArray, 0, bigArray.length));2.3 ForkJoinPool vs ThreadPoolExecutor
| 特性 | ThreadPoolExecutor | ForkJoinPool |
|---|---|---|
| 队列 | 共享队列 | 每线程独立队列 + 工作窃取 |
| 适用场景 | IO 密集型、独立任务 | CPU 密集型、可拆分任务 |
| 线程数 | 静态/手动调整 | 动态调整 |
| 典型应用 | Web 请求处理 | 并行排序、大数据聚合 |
CompletableFuture 默认使用 ForkJoinPool,因为异步任务链本身就是"分拆 -> 合并"的模式,和 ForkJoinPool 天然匹配。
三、并发容器
3.1 ConcurrentHashMap
JDK 8 的 ConcurrentHashMap 使用 CAS + synchronized 保证线程安全:
- 数组中某个位置为空:CAS 直接放入。
- 数组中某个位置已有数据(hash 冲突):synchronized 锁住该位置的头节点,在链表/红黑树上操作。
结构和 HashMap 一样是 数组 + 链表 + 红黑树。链表长度 >= 8 且数组长度 >= 64 时转为红黑树。
关键区别:ConcurrentHashMap 不允许 key 或 value 为 null。因为在并发场景下无法区分"这个 key 对应的 value 就是 null"和"这个 key 不存在"。
3.2 CopyOnWriteArrayList
写时复制的线程安全 ArrayList:每次写操作都会创建底层数组的副本,在副本上修改,修改完后替换原数组。读操作不加锁,直接读原数组。
CopyOnWriteArrayList<String> list = new CopyOnWriteArrayList<>();
list.add("hello"); // 加锁 → 复制数组 → 在新数组末尾添加 → 替换原数组
String s = list.get(0); // 无锁直接读适用场景:读多写少(如监听器列表、配置缓存)。写多或数据量大时慎用,因为每次写都要拷贝整个数组。
弱一致性:写操作完成前,读操作看到的还是旧数组。
3.3 ConcurrentLinkedQueue
基于 CAS 实现的无界非阻塞队列,适合高并发场景下的"生产-消费"模式。它不支持阻塞等待,poll() 拿不到就返回 null。
四、阻塞队列
阻塞队列是线程池的核心组件之一,也是生产者-消费者模式的标准实现。核心行为:队列空时 take() 阻塞,队列满时 put() 阻塞。
4.1 七大阻塞队列对比
| 队列 | 底层 | 是否有界 | 特点 |
|---|---|---|---|
| ArrayBlockingQueue | 数组 | 有界 | 最常用,创建时必须指定容量 |
| LinkedBlockingQueue | 链表 | 可有界可无界 | 默认 Integer.MAX_VALUE(注意 OOM) |
| PriorityBlockingQueue | 堆 | 无界 | 按优先级出队 |
| DelayQueue | 堆 | 无界 | 元素到期才能取出 |
| SynchronousQueue | 无 | 容量为 0 | 每个 put 必须等到对应的 take |
| LinkedTransferQueue | 链表 | 无界 | 支持直接传递元素 |
| LinkedBlockingDeque | 双向链表 | 可有界 | 双端队列,两端都可存取 |
4.2 生产者-消费者模式
BlockingQueue<String> queue = new ArrayBlockingQueue<>(100);
// 生产者
new Thread(() -> {
while (true) {
queue.put(produce()); // 队列满时自动阻塞
}
}).start();
// 消费者
new Thread(() -> {
while (true) {
String item = queue.take(); // 队列空时自动阻塞
consume(item);
}
}).start();这个模式比手动 wait()/notify() 简洁得多,而且线程安全。
五、异步编程
5.1 FutureTask
FutureTask 是 Future 接口的基本实现,结合 Callable 使用,可以获取异步任务的返回值。
FutureTask<String> task = new FutureTask<>(() -> {
Thread.sleep(2000);
return "done";
});
new Thread(task).start();
String result = task.get(3, TimeUnit.SECONDS); // 阻塞等待,最多3秒FutureTask 内部维护了一个 state 状态机(NEW → COMPLETING → NORMAL/EXCEPTIONAL/CANCELLED),以及一个 WaitNode 链表来管理等待结果的线程。get() 时如果任务未完成,当前线程会被加入链表并通过 LockSupport.park() 挂起,任务完成后逐个 unpark() 唤醒。
5.2 CompletableFuture
JDK 8 引入的 CompletableFuture 是异步编程的利器,支持链式调用和任务编排,比 FutureTask 强大得多。
核心三个函数式接口:
| 接口 | 入参 | 返回值 | 用途 |
|---|---|---|---|
| Supplier<U> | 无 | 有 | 生产者 |
| Consumer<T> | 有 | 无 | 消费者 |
| Function<T,U> | 有 | 有 | 转换器 |
常用 API:
// 无返回值的异步任务
CompletableFuture.runAsync(() -> doTask());
// 有返回值的异步任务
CompletableFuture<String> future = CompletableFuture.supplyAsync(() -> "hello");
// 串行:A 完了 B 接着干,B 需要 A 的结果
future.thenApply(result -> result.toUpperCase());
// 串行:B 消费 A 的结果,不返回
future.thenAccept(result -> log(result));
// 串行:B 不需要 A 的结果
future.thenRun(() -> cleanUp());
// 并行:A 和 B 都完成后再执行 C
futureA.thenCombine(futureB, (a, b) -> a + b);
// 并行:等所有完成
CompletableFuture.allOf(f1, f2, f3).thenRun(() -> allDone());
// 并行:等任意一个完成
CompletableFuture.anyOf(f1, f2, f3).thenAccept(first -> useFirst(first));
// 异常处理
future.exceptionally(ex -> fallback());
future.handle((result, ex) -> ex != null ? fallback() : result);注意:CompletableFuture 如果不指定线程池,默认使用 ForkJoinPool.commonPool(),其中的线程是守护线程。如果 main 方法结束了,守护线程也会跟着结束。生产环境建议显式传入自定义线程池。
带 Async 后缀和不带的区别:thenApply 在前一个任务的线程上执行,thenApplyAsync 提交到线程池异步执行。
六、并发工具类总结
| 工具 | 核心用途 | 底层 |
|---|---|---|
| CountDownLatch | N 个任务都完成后再继续 | AQS 共享模式 |
| CyclicBarrier | N 个线程互相等待到齐后一起出发 | ReentrantLock + Condition |
| Semaphore | 控制同时访问资源的线程数 | AQS 共享模式 |
| Phaser | 分阶段同步(CountDownLatch + CyclicBarrier 的超集) | 独立实现 |
三者的使用场景对比:
- CountDownLatch:主线程等多个工作线程完成,像"火箭倒计时"。
- CyclicBarrier:多个线程互相等待,像"旅行团集合"。
- Semaphore:限制并发数,像"停车场限流"。
七、常见面试题精选
Q1:ForkJoinPool 和 ThreadPoolExecutor 的区别?
ThreadPoolExecutor 基于任务分配,所有线程共享一个任务队列。ForkJoinPool 基于工作窃取,每个线程有自己的双端队列,空闲线程从其他线程队列尾部偷任务。ForkJoinPool 适合 CPU 密集型的可拆分任务(如并行排序),ThreadPoolExecutor 适合 IO 密集型的独立任务(如 Web 请求处理)。CompletableFuture 默认用 ForkJoinPool。
Q2:为什么不建议用 Executors 创建线程池?
newFixedThreadPool 和 newSingleThreadExecutor 使用无界的 LinkedBlockingQueue,任务堆积可能导致 OOM。newCachedThreadPool 的最大线程数是 Integer.MAX_VALUE,可能创建大量线程导致 OOM。应该用 ThreadPoolExecutor 自行构造,给队列指定容量。
Q3:线程数设定成多少合适?
没有万能公式。CPU 密集型参考 N+1,IO 密集型参考 2N+1(N 为 CPU 核数)。更精确的公式是 核数 x 目标利用率 x (1 + 等待时间/计算时间)。但最靠谱的做法是先给初始值,然后通过压测不断调整。
Q4:线程池有哪些核心参数?
7 个:corePoolSize(核心线程数)、maximumPoolSize(最大线程数)、keepAliveTime(非核心线程空闲存活时间)、unit(时间单位)、workQueue(任务队列)、threadFactory(线程工厂)、handler(拒绝策略)。
Q5:线程池的拒绝策略有哪些?
四种内置策略:AbortPolicy(抛异常,默认)、DiscardPolicy(默默丢弃)、DiscardOldestPolicy(丢最老的任务)、CallerRunsPolicy(调用者线程自己执行)。生产环境最常用 AbortPolicy 和 CallerRunsPolicy。
Q6:介绍下 JUC,有哪些工具类?
JUC(java.util.concurrent)是 Java 并发编程的核心包。主要分四类:1)线程池与任务执行(ThreadPoolExecutor、ForkJoinPool、CompletableFuture);2)同步器(ReentrantLock、Semaphore、CountDownLatch、CyclicBarrier,底层都是 AQS);3)并发集合(ConcurrentHashMap、CopyOnWriteArrayList、各种 BlockingQueue);4)原子变量(AtomicInteger、LongAdder 等,基于 CAS)。
小结
并发编程全景图
┌──────────────────────────────────────────────────┐
│ 线程基础 │ synchronized │ volatile │ JMM │
├──────────────────────────────────────────────────┤
│ AQS (ReentrantLock / Semaphore / ...) │
├──────────────────────────────────────────────────┤
│ 原子类 │ ThreadLocal │ StampedLock │
├──────────────────────────────────────────────────┤
│ 线程池 (ThreadPoolExecutor / ForkJoinPool) │
├──────────────────────────────────────────────────┤
│ 并发容器 (ConcurrentHashMap / CopyOnWrite / ...)│
├──────────────────────────────────────────────────┤
│ 阻塞队列 (ArrayBQ / LinkedBQ / DelayQueue / ...)│
├──────────────────────────────────────────────────┤
│ 异步编程 (FutureTask / CompletableFuture) │
└──────────────────────────────────────────────────┘并发编程归根到底就是解决三件事:
- 互斥:同一时间只有一个线程能操作共享资源(锁、CAS)。
- 同步:线程之间协调执行顺序(等待/通知、CountDownLatch、CyclicBarrier)。
- 通信:线程之间传递数据(阻塞队列、Future、共享变量 + volatile)。
掌握了这三件事对应的工具和原理,并发编程就不再是拦路虎了。
附录:核心组件源码要点
A.1 ThreadPoolExecutor 的 execute 方法
public void execute(Runnable command) {
if (command == null) throw new NullPointerException();
int c = ctl.get();
// 1. 工作线程数 < 核心线程数 → 创建核心线程
if (workerCountOf(c) < corePoolSize) {
if (addWorker(command, true)) return;
c = ctl.get();
}
// 2. 线程池运行中 且 入队成功
if (isRunning(c) && workQueue.offer(command)) {
int recheck = ctl.get();
// 再次检查状态,可能刚好 shutdown 了
if (!isRunning(recheck) && remove(command))
reject(command);
// 可能核心线程数为 0 或者允许超时导致没有工作线程
else if (workerCountOf(recheck) == 0)
addWorker(null, false);
}
// 3. 入队失败 → 创建非核心线程
else if (!addWorker(command, false))
reject(command); // 非核心也满了 → 拒绝
}A.2 Worker 的工作循环
Worker 的核心是 runWorker 方法:
final void runWorker(Worker w) {
Runnable task = w.firstTask;
w.firstTask = null;
while (task != null || (task = getTask()) != null) {
w.lock();
try {
beforeExecute(w.thread, task); // 钩子方法
task.run();
afterExecute(task, null); // 钩子方法
} finally {
task = null;
w.completedTasks++;
w.unlock();
}
}
// 跳出循环意味着 getTask() 返回 null
processWorkerExit(w, false);
}Worker 不只执行 firstTask,而是持续从队列获取任务。getTask() 返回 null 才会退出,退出意味着这个线程要被销毁了。
A.3 getTask 方法的退出时机
private Runnable getTask() {
boolean timedOut = false;
for (;;) {
int c = ctl.get();
int rs = runStateOf(c);
// 线程池 STOP 或 (SHUTDOWN 且队列空) → 返回 null
if (rs >= SHUTDOWN && (rs >= STOP || workQueue.isEmpty())) {
decrementWorkerCount();
return null;
}
int wc = workerCountOf(c);
boolean timed = allowCoreThreadTimeOut || wc > corePoolSize;
// 工作线程多于核心线程 且 上次取任务超时 → 返回 null(该线程将被销毁)
if ((wc > maximumPoolSize || (timed && timedOut))
&& (wc > 1 || workQueue.isEmpty())) {
if (compareAndDecrementWorkerCount(c))
return null;
continue;
}
try {
Runnable r = timed ?
workQueue.poll(keepAliveTime, TimeUnit.NANOSECONDS) :
workQueue.take(); // 核心线程死等
if (r != null) return r;
timedOut = true; // poll 超时
} catch (InterruptedException retry) {
timedOut = false;
}
}
}关键点:核心线程调用 take() 死等,非核心线程调用 poll(keepAliveTime) 限时等待。超时后 timedOut 变为 true,下一轮循环就会返回 null。
A.4 ConcurrentHashMap 的 putVal 要点
JDK 8 的 ConcurrentHashMap 在 putVal 中使用了三种策略:
final V putVal(K key, V value, boolean onlyIfAbsent) {
if (key == null || value == null) throw new NullPointerException();
int hash = spread(key.hashCode());
for (Node<K,V>[] tab = table;;) {
Node<K,V> f; int n, i, fh;
if (tab == null || (n = tab.length) == 0)
tab = initTable(); // 懒初始化
else if ((f = tabAt(tab, i = (n - 1) & hash)) == null) {
if (casTabAt(tab, i, null, new Node<>(hash, key, value, null)))
break; // CAS 直接放入
}
else if ((fh = f.hash) == MOVED)
tab = helpTransfer(tab, f); // 协助扩容
else {
synchronized (f) { // 锁住头节点
// 链表或红黑树操作
}
}
}
addCount(1L, binCount);
return null;
}三种策略:
- 数组位置为空:CAS 放入,无锁。
- 正在扩容:当前线程协助扩容。
- 有 hash 冲突:synchronized 锁住该位置的头节点,在链表/红黑树上操作。
spread() 方法将 hashCode 的高 16 位参与运算,并保证结果为正数(因为负数 hash 有特殊含义:-1 表示正在扩容,-2 表示红黑树节点)。
A.5 ConcurrentHashMap 不允许 null 的原因
HashMap 允许 key 和 value 为 null,但 ConcurrentHashMap 不允许。原因是:在并发环境下,get(key) 返回 null 时,你无法区分"key 存在但 value 是 null"还是"key 不存在"。在单线程 HashMap 中可以通过 containsKey() 再次确认,但在并发场景下两次调用之间可能有其他线程修改了 Map,导致 containsKey() 的结果和 get() 的结果不一致。
A.6 CompletableFuture 的任务编排原理
CompletableFuture 内部维护了一个 stack 栈结构(Treiber Stack),用来存放后置任务。
当你调用 thenApply(fn) 时:
- 先检查前置任务是否已完成(result != null)。
- 如果已完成,直接执行 fn。
- 如果未完成,将 fn 封装为 UniApply 节点,push 到 stack 中。
- push 之后还会再检查一次(防止 push 期间前置任务刚好完成),确保任务不会遗漏。
前置任务完成时,会调用 postComplete() 方法,从 stack 中依次弹出后置任务并执行。如果后置任务自己也有后置任务(嵌套),会递归处理。
A.7 FutureTask 的状态机
FutureTask 使用 volatile int state 管理任务状态:
NEW(0) → COMPLETING(1) → NORMAL(2) // 正常完成
NEW(0) → COMPLETING(1) → EXCEPTIONAL(3) // 异常完成
NEW(0) → CANCELLED(4) // 取消(不中断)
NEW(0) → INTERRUPTING(5) → INTERRUPTED(6) // 取消(中断)get() 时如果 state <= COMPLETING,当前线程封装为 WaitNode 加入等待链表并 park。任务完成后 finishCompletion() 遍历链表逐个 unpark。
A.8 并发编程三大特性速查
| 特性 | 含义 | 保证方式 |
|---|---|---|
| 原子性 | 操作不可分割 | synchronized / CAS / Lock |
| 可见性 | 修改对其他线程立即可见 | volatile / synchronized / Lock / final |
| 有序性 | 禁止指令重排 | volatile(内存屏障)/ happens-before 规则 |
volatile 保证可见性和有序性,但不保证原子性(如 i++ 不是原子操作)。
synchronized 三者全保证,但性能开销最大。
CAS 保证原子性,配合 volatile 保证可见性。
A.9 JMM 与 happens-before
Java 内存模型(JMM)通过 happens-before 规则定义了多线程间操作的可见性保证:
- 程序顺序规则:同一线程中,前面的操作 happens-before 后面的操作。
- Monitor 锁规则:unlock happens-before 后续的 lock。
- volatile 规则:写 volatile happens-before 后续的读 volatile。
- 线程启动规则:start() happens-before 子线程的第一个操作。
- 线程终止规则:子线程的最后一个操作 happens-before join() 返回。
- 传递性:A happens-before B,B happens-before C,则 A happens-before C。
A.10 读写锁 ReentrantReadWriteLock
ReentrantReadWriteLock 把 AQS 的 state 拆成两半用:
- 高 16 位:读锁的持有数(共享锁)。
- 低 16 位:写锁的重入数(独占锁)。
读读不互斥,读写互斥,写写互斥。每个读线程用 ThreadLocal 记录自己的重入次数。
写锁饥饿问题:如果读锁一直被持有,写线程就永远拿不到写锁。解决方式是:当有写线程在队列中等待时,后续的读线程不能"插队"直接获取读锁,必须排在写线程后面。
A.11 线程池的钩子方法
ThreadPoolExecutor 提供了三个钩子方法供子类扩展:
protected void beforeExecute(Thread t, Runnable r) { }
protected void afterExecute(Runnable r, Throwable t) { }
protected void terminated() { }beforeExecute 和 afterExecute 在每个任务执行前后被调用,可以用来做日志记录、性能统计、线程上下文清理等。terminated 在线程池完全终止后调用。
class MonitoringPool extends ThreadPoolExecutor {
@Override
protected void beforeExecute(Thread t, Runnable r) {
log.info("任务开始: {}", r);
}
@Override
protected void afterExecute(Runnable r, Throwable t) {
if (t != null) log.error("任务异常", t);
log.info("任务结束: {}", r);
}
}A.12 Semaphore 的 acquire 流程
Semaphore 的 acquire() 使用 AQS 共享模式:
tryAcquireShared(permits):获取当前 state,减去需要的许可数。- 如果结果 >= 0,CAS 更新 state,获取成功。
- 如果结果 < 0,返回负数表示获取失败。
- 获取失败的线程封装为 Node(共享模式),加入同步队列并 park。
release()时 CAS 增加 state,然后唤醒队列中的等待线程。
公平版在 tryAcquireShared 中会先检查队列(hasQueuedPredecessors()),非公平版直接 CAS 抢。
A.13 CyclicBarrier 的源码要点
CyclicBarrier 不是基于 AQS 实现的,而是基于 ReentrantLock + Condition:
public class CyclicBarrier {
private final ReentrantLock lock = new ReentrantLock();
private final Condition trip = lock.newCondition();
private final int parties; // 总参与方数
private final Runnable barrierCommand; // 到齐后执行的回调
private int count; // 剩余等待数
private Generation generation; // 代次标记
}await() 的核心逻辑:
- 加锁。
count--。- 如果 count == 0(所有人到齐了):
- 执行 barrierCommand(如果有)。
- 调用
nextGeneration():唤醒所有等待线程,重置 count,创建新的 generation。
- 如果 count > 0(还有人没到):
- 调用
trip.await()挂起当前线程。
- 调用
CyclicBarrier 可重用的秘密就在于 nextGeneration():它重置了 count 并创建新的 generation 对象,所以下一轮可以继续使用。
A.14 CountDownLatch 的源码要点
CountDownLatch 基于 AQS 共享模式:
- 构造:
setState(count),将计数值设为 AQS 的 state。 - countDown():CAS 将 state 减 1,减到 0 时调用
doReleaseShared()唤醒所有等待线程。 - await():检查 state 是否为 0,为 0 直接返回;否则封装为共享 Node 加入队列并 park。
// CountDownLatch 的 tryAcquireShared
protected int tryAcquireShared(int acquires) {
return (getState() == 0) ? 1 : -1;
}
// CountDownLatch 的 tryReleaseShared
protected boolean tryReleaseShared(int releases) {
for (;;) {
int c = getState();
if (c == 0) return false;
int nextc = c - 1;
if (compareAndSetState(c, nextc))
return nextc == 0; // 减到 0 才返回 true,触发唤醒
}
}A.15 DelayQueue 的应用场景
DelayQueue 是一个无界阻塞队列,元素必须实现 Delayed 接口,只有到期的元素才能被取出。常见应用场景:
- 重试机制:接口调用失败后,将请求放入 DelayQueue,延迟一段时间后重试。可以用指数退避(2s、4s、8s、16s......)。
- 订单超时取消:订单创建时放入 DelayQueue,30 分钟后如果还未支付就取出并取消。
- 缓存过期:缓存条目放入 DelayQueue,到期后取出并清除。
class DelayedTask implements Delayed {
private long deadline;
@Override
public long getDelay(TimeUnit unit) {
return unit.convert(deadline - System.nanoTime(), TimeUnit.NANOSECONDS);
}
@Override
public int compareTo(Delayed other) {
return Long.compare(this.deadline, ((DelayedTask)other).deadline);
}
}A.16 SynchronousQueue 的特殊性
SynchronousQueue 是一个容量为 0 的队列,每个 put() 操作必须等到对应的 take() 操作,反之亦然。它不存储任何元素,而是直接在两个线程之间传递数据。
newCachedThreadPool 用的就是 SynchronousQueue:任务提交后如果有空闲线程就直接交给它,没有就创建新线程。这也解释了为什么 CachedThreadPool 的核心线程数是 0:它不需要任务排队,要么立即执行,要么创建线程执行。