Java 从入门到精通(十三):JUC 核心工具——AQS、Lock、原子类与并发容器

本文默认读者已掌握线程生命周期、JMM 与 synchronized/volatile(系列第十二篇),JDK 以 8 为主,涉及 9 之后的变化(Unsafe 迁往 VarHandle)会单独标注。AQS 一节最硬,建议按”state 是什么 → 队列怎么排 → 线程怎么睡 → 谁把它叫醒”四步推演,读完再回头看 ReentrantLock、CountDownLatch、Semaphore,会发现它们只是同一套骨架上的不同血肉。

一、JUC 全景与定位

靠 synchronized 与 volatile 能解决大部分并发问题,但它们有两个天然短板:一是不可中断——等监视器锁时无法响应 interrupt(),死锁只能重启;二是不够灵活——没有尝试加锁、没有超时、没有读写分离、也没有多个等待队列,想做”生产者只唤醒消费者”这种精细控制,wait/notify 只能 notifyAll 全量广播。

JUC(java.util.concurrent)正是为了补上这块拼图。它由 Doug Lea 主导设计,从 JDK 5 引入,核心思路是:把并发控制的公共骨架抽出来做成可复用组件,让业务代码从”手工操作 wait/notify”升级到”组合现成工具”。

1.1 包的四层结构

子包/包 代表类 解决的核心问题 本篇覆盖
java.util.concurrent.locks AbstractQueuedSynchronizer、ReentrantLock、ReentrantReadWriteLock、StampedLock 显式锁、条件队列、读写分离、乐观读 二、三章
java.util.concurrent.atomic AtomicInteger、AtomicReference、LongAdder、AtomicStampedReference 无锁原子更新、累加器 四章
java.util.concurrent(集合部分) ConcurrentHashMap、CopyOnWriteArrayList、ConcurrentLinkedQueue、各类 BlockingQueue 线程安全容器的性能与语义 五章
java.util.concurrent(协作部分) CountDownLatch、CyclicBarrier、Semaphore、Phaser、Exchanger 线程间的等待、汇合、限流 六章
java.util.concurrent(执行器) ThreadPoolExecutor、FutureTask、CompletableFuture、ForkJoinPool 任务调度与异步编排 下一篇
附带 ThreadLocal、ThreadLocalRandom 线程封闭、上下文传递 七章

1.2 与 synchronized 的分工

很多人纠结”有了 synchronized 为什么还要 Lock”。答案不是替换,而是分工:synchronized 是 JVM 内置的 Monitor 机制,解锁由编译器插入的 monitorenter/monitorexit 保证,JVM 还会做锁粗化、锁消除与自适应自旋,JDK 6 之后性能已持平,能用就用;Lock 是纯 Java 实现的显式锁,胜在能力——可中断、可超时、可非阻塞尝试、可选公平性、一把锁绑多个 Condition、支持读写分离与乐观读。需要其中一项时再切换。

选择顺序建议:先想能不能”不加锁”(无锁结构、线程封闭、不可变对象)→ 再想 synchronized 够不够 → 需要高级能力才上 ReentrantLock → 读多写极少考虑 StampedLock 或 CopyOnWrite → 高并发计数优先 LongAdder。

1.3 本篇与下一篇的边界

关注点 本篇(十三) 下一篇(十四)
核心主题 同步原语与容器:如何”安全地共享数据” 任务执行框架:如何”高效地执行任务”
关键抽象 AQS、Lock、Condition、原子类、BlockingQueue Executor、Future、CompletableFuture、ForkJoinPool
典型问题 竞态、可见性、死锁、伪共享、内存泄漏 线程池参数、拒绝策略、异步编排、任务窃取
交集 BlockingQueue 既是容器也是线程池的工作队列;FutureTask 内部也用到了 AQS 的共享模式,本篇会顺带点出 线程池的 worker 抢任务机制会复用本篇的 CAS 知识

二、AQS 原理:一把锁的骨架

AbstractQueuedSynchronizer(下文简称 AQS)是 JUC 的心脏。ReentrantLock、Semaphore、CountDownLatch、ReentrantReadWriteLock、FutureTask、ThreadPoolExecutor.Worker 全部直接或间接继承自它。理解 AQS,等于一次性理解了半个 JUC。

2.1 模板方法模式

AQS 用了非常经典的模板方法模式:把排队、阻塞、唤醒、取消这些与业务无关的流程全部实现在基类,只把”能不能拿到资源”这一个判断留给子类。子类需重写的方法只有五个,全是 protected:

方法 模式 语义
tryAcquire(int) 独占 尝试获取资源,成功返回 true
tryRelease(int) 独占 尝试释放资源
tryAcquireShared(int) 共享 返回负数失败,0 成功但无剩余,正数成功且有剩余
tryReleaseShared(int) 共享 释放共享资源
isHeldExclusively() 独占 当前线程是否独占持有,供 Condition 使用

其余 acquire、acquireInterruptibly、acquireShared、release 等全是 final 模板方法,子类不可改。这就是 AQS 的优雅之处:流程固定,语义可插拔。

2.2 核心三件套

2.2.1 volatile int state

state 是同步器的资源计数器,volatile 保证可见性,配 getState/setState/compareAndSetState 访问(CAS 由 Unsafe 提供)。其含义由子类定义:

  • ReentrantLock:state=0 表示未锁定,>0 表示被持有,数值即重入次数;
  • Semaphore:state 表示剩余许可数;
  • CountDownLatch:state 表示还没完成的计数;
  • ReentrantReadWriteLock:state 高 16 位存读锁数量,低 16 位存写锁重入次数——一个 int 拆成两个半字用。

2.2.2 CLH 队列的变体

AQS 内部维护一条 FIFO 双向队列,官方注释称之为 “variant of CLH queue”。相对原始 CLH(单向、靠忙等前驱状态位)做了三点改造:

  1. 改双向:Node 同时持有 prev 与 next。单向链表在节点 CANCELLED 时拿不到前驱,无法从中间摘除;双向可直接 prev.next = next; next.prev = prev。
  2. 自旋改 park:竞争激烈、临界区长时忙等纯属烧 CPU,AQS 改为自旋失败后 LockSupport.park() 挂起,由前驱释放时 unpark。
  3. 引入 head 哑节点:head 指向”当前已持有资源的线程”或空壳节点,真正排队的是 head.next。

不用普通 LinkedList 的原因是无锁入队:enq 用 CAS 抢设 tail,失败自旋重试,全程不加锁;而 prev 指针让”前驱取消就往前跳”实现得非常干净。

2.2.3 Node 与 waitStatus

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
static final class Node {
// 标记:共享模式 / 独占模式
static final Node SHARED = new Node();
static final Node EXCLUSIVE = null;

// waitStatus 的四种取值(0 是第三种:新建或已完成)
static final int CANCELLED = 1; // 等待超时或被中断,需要移出队列
static final int SIGNAL = -1; // 后继节点在等我释放/取消,我必须唤醒它
static final int CONDITION = -2; // 节点正躺在 Condition 的等待队列里
static final int PROPAGATE = -3; // 共享模式下,唤醒需要向后持续传播

volatile int waitStatus;
volatile Node prev;
volatile Node next;
volatile Thread thread; // 节点代表的线程
Node nextWaiter; // Condition 队列的单向链,或 SHARED 标记
}

waitStatus 是理解 AQS 的钥匙,五个取值必须记牢:

  • CANCELLED = 1:唯一正数。节点因超时或中断放弃竞争,一旦置位不再变化,会在 cancelAcquire 中被摘链。
  • SIGNAL = -1:最核心。表示”我的后继在 park,我有义务释放时 unpark 它”。所以每个节点入队后必须先确保前驱是 SIGNAL 才敢 park,否则会”睡着了没人叫”。
  • CONDITION = -2:节点不在同步队列,而在某个 ConditionObject 的单向等待队列里,signal 时搬回同步队列并把状态改回 0。
  • PROPAGATE = -3:共享模式专用。多个线程可同时持有资源,一次 releaseShared 可能要连续唤醒多个后继,该状态保证唤醒能向后传播,避免丢失。
  • 0:新建节点的默认值,表示”当前无事发生”。

2.3 acquire 全流程推演

独占模式获取资源的入口是 acquire(int arg),逻辑只有三行,但背后是完整的一套排队机制:

1
2
3
4
5
public final void acquire(int arg) {
if (!tryAcquire(arg) && // ① 先抢一次
acquireQueued(addWaiter(Node.EXCLUSIVE), arg)) // ② 抢不到:入队 → 排队等待
selfInterrupt(); // ③ 等待期间被中断过,补上中断标记
}

完整流程分五步:

  1. 快速尝试:调子类 tryAcquire。非公平锁下这步常直接成功(刚释放的锁被新线程抢走),这是非公平吞吐更高的根源。
  2. 包装入队:addWaiter 把当前线程包成 EXCLUSIVE 节点,先 CAS 一次挂队尾,失败则进 enq 自旋 CAS 直到成功。
  3. 排队自旋:acquireQueued 是 for(;;):前驱是 head 就再 tryAcquire,失败则问 shouldParkAfterFailedAcquire 能否 park。
  4. 安全入睡:shouldParkAfterFailedAcquire 先把前驱 waitStatus CAS 成 SIGNAL 并返回 false 让外层再转一圈;第二圈发现前驱已是 SIGNAL,才真正 LockSupport.park(this)。
  5. 被唤醒重来:前驱释放时 unparkSuccessor 找到 head 后第一个非 CANCELLED 节点并 unpark。被唤醒线程回到第 3 步,tryAcquire 成功后自己成为新 head 并返回中断标记。
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
final boolean acquireQueued(final Node node, int arg) {
boolean failed = true;
try {
boolean interrupted = false;
for (;;) {
final Node p = node.predecessor();
// 只有前驱是 head 才尝试获取:保证 FIFO 公平,避免队尾插队
if (p == head && tryAcquire(arg)) {
setHead(node); // 自己成为 head,thread 置 null
p.next = null; // 帮助 GC 回收旧 head
failed = false;
return interrupted;
}
// 判断是否需要 park;返回 true 就 park,醒来后检查中断位
if (shouldParkAfterFailedAcquire(p, node) &&
parkAndCheckInterrupt())
interrupted = true; // 注意:不抛异常,只记账,最后补中断
}
} finally {
if (failed) cancelAcquire(node); // 异常才取消排队
}
}

private static boolean shouldParkAfterFailedAcquire(Node pred, Node node) {
int ws = pred.waitStatus;
if (ws == Node.SIGNAL)
return true; // 前驱承诺会唤醒我,可以放心睡
if (ws > 0) {
// 前驱已取消,向前跳过所有 CANCELLED 节点,再修正链表
do {
node.prev = pred = pred.prev;
} while (pred.waitStatus > 0);
pred.next = node;
} else {
// 前驱是 0 或 PROPAGATE:用 CAS 把它改成 SIGNAL,下次循环再 park
compareAndSetWaitStatus(pred, ws, Node.SIGNAL);
}
return false;
}

容易忽略的细节:acquire 不响应中断。parkAndCheckInterrupt() 只把中断标记记下来继续排队,等真正拿到资源后才 selfInterrupt() 补上;想响应中断要用 acquireInterruptibly,它检测到中断直接抛 InterruptedException 并 cancelAcquire。

2.4 独占模式与共享模式的差异

对比项 独占模式 EXCLUSIVE 共享模式 SHARED
获取方法 acquire / tryAcquire 返回 boolean acquireShared / tryAcquireShared 返回 int
成功语义 只有当前线程能持有 多个线程可同时持有
节点标记 nextWaiter == null nextWaiter == SHARED
唤醒行为 release 只唤醒 head 的一个后继 releaseShared 唤醒后可能继续向后传播(PROPAGATE)
典型实现 ReentrantLock Semaphore、CountDownLatch、ReadLock

共享模式的关键在 setHeadAndPropagate:新 head 就位后,若还有剩余资源(tryAcquireShared 返回 > 0)或 head 状态为 PROPAGATE,会继续 doReleaseShared 唤醒下一个节点,形成级联唤醒。CountDownLatch 计数归零时,所有 await 线程正是靠它一次性全部放行。

2.5 各类同步器如何复用 AQS

同步器 state 的含义 模式 关键 tryAcquire 逻辑
ReentrantLock 0 未锁;N 表示重入 N 次 独占 CAS 改 0→1;若已是自己则 state+1
ReentrantReadWriteLock 高 16 位读计数,低 16 位写重入 读共享/写独占 读锁看写锁是否被占;写锁看 state 是否非 0
Semaphore 剩余许可数 共享 自旋 CAS 把 state 减掉 acquire 的许可数,不够则返回负数
CountDownLatch 未完成的任务计数 共享 只有 state == 0 才返回 1,否则 -1
FutureTask 任务状态(NEW/COMPLETING 等) 共享 任务完成才返回 1

2.6 手写一个不可重入独占锁

看懂上面之后,自己实现一个锁只需四十行。tryAcquire 里”已锁定就直接返回 false”这一句,正是”不可重入”的定义。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
import java.util.concurrent.TimeUnit;
import java.util.concurrent.locks.AbstractQueuedSynchronizer;
import java.util.concurrent.locks.Condition;

/**
* 基于 AQS 的不可重入独占锁(Mutex)
* 约定:state == 0 表示未锁定,state == 1 表示已锁定
*/
public class MutexLock implements java.io.Serializable {
private static final long serialVersionUID = 1L;

/** 静态内部类:只重写 tryAcquire / tryRelease / isHeldExclusively */
private static class Sync extends AbstractQueuedSynchronizer {
// 是否被占用
@Override
protected boolean isHeldExclusively() {
return getState() == 1;
}

// state 为 0 时用 CAS 抢成 1,抢到即获得锁;不可重入:已占用直接失败
@Override
protected boolean tryAcquire(int acquires) {
assert acquires == 1;
if (compareAndSetState(0, 1)) {
setExclusiveOwnerThread(Thread.currentThread()); // 记录持有线程
return true;
}
return false;
}

// 释放:state 置 0,必须先把持有者清空再 setState,保证可见顺序
@Override
protected boolean tryRelease(int releases) {
assert releases == 1;
if (getState() == 0) throw new IllegalMonitorStateException();
setExclusiveOwnerThread(null);
setState(0);
return true;
}

Condition newCondition() { return new ConditionObject(); }
}

private final Sync sync = new Sync();

public void lock() { sync.acquire(1); }
public void unlock() { sync.release(1); }
public boolean tryLock() { return sync.tryAcquire(1); }
public boolean isLocked() { return sync.isHeldExclusively(); }

/** 可中断地获取锁,等待期间被 interrupt 会抛异常 */
public void lockInterruptibly() throws InterruptedException {
sync.acquireInterruptibly(1);
}

/** 带超时的尝试,arg 为 1,超时时间由 AQS 的 nanosTimeout 处理 */
public boolean tryLock(long timeout, TimeUnit unit) throws InterruptedException {
return sync.tryAcquireNanos(1, unit.toNanos(timeout));
}

public Condition newCondition() { return sync.newCondition(); }

// ---- 测试:10 个线程各加 1000 次,结果应为 10000 ----
public static void main(String[] args) throws InterruptedException {
MutexLock lock = new MutexLock();
int[] counter = {0};
Thread[] ts = new Thread[10];
for (int i = 0; i < ts.length; i++) {
ts[i] = new Thread(() -> {
for (int j = 0; j < 1000; j++) {
lock.lock();
try {
counter[0]++; // 临界区:非原子操作
} finally {
lock.unlock(); // 必须放 finally,防止异常导致锁泄露
}
}
});
ts[i].start();
}
for (Thread t : ts) t.join();
System.out.println("counter = " + counter[0]); // 稳定输出 10000
}
}

把它改成可重入只需加一句判断:若 getExclusiveOwnerThread() == Thread.currentThread(),则 setState(getState() + 1) 返回 true,释放时对应减一、减到 0 才真正置 0——这正是 ReentrantLock 的实现。

三、Lock 体系:从接口到实现类

3.1 Lock 接口的方法语义

1
2
3
4
5
6
7
8
public interface Lock {
void lock(); // 拿不到就阻塞,不可中断
void lockInterruptibly() throws InterruptedException; // 阻塞但响应中断
boolean tryLock(); // 立即返回成败,不阻塞
boolean tryLock(long time, TimeUnit unit) throws InterruptedException; // 限时尝试
void unlock(); // 必须释放,否则死锁
Condition newCondition(); // 绑定一个等待队列
}

对比 synchronized 缺失的能力:tryLock 可做死锁规避(拿不到第二把锁就先放开第一把,退避重试);lockInterruptibly 让死锁可被外部解除;newCondition 支持一把锁开多个等待队列。

3.2 ReentrantLock 的可重入实现

可重入指”同一线程可多次获取同一把锁”,避免自己锁自己。ReentrantLock 靠两点实现:state 计数与持有线程判定:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
// ReentrantLock.NonfairSync.tryAcquire(JDK 8 精简版)
final boolean nonfairTryAcquire(int acquires) {
final Thread current = Thread.currentThread();
int c = getState();
if (c == 0) {
// 无锁:CAS 抢一下,抢到就把持有者设为自己
if (compareAndSetState(0, acquires)) {
setExclusiveOwnerThread(current);
return true;
}
} else if (current == getExclusiveOwnerThread()) {
// 已锁且是自己:重入,state 累加(int 溢出会抛 Error,故上限极大)
int nextc = c + acquires;
if (nextc < 0) throw new Error("Maximum lock count exceeded");
setState(nextc); // 只有自己能改,无需 CAS
return true;
}
return false; // 锁被别人持有
}

protected final boolean tryRelease(int releases) {
int c = getState() - releases;
if (Thread.currentThread() != getExclusiveOwnerThread())
throw new IllegalMonitorStateException(); // 未持有却释放,直接报错
boolean free = false;
if (c == 0) { // 完全释放:次数归零才算真正解锁
free = true;
setExclusiveOwnerThread(null);
}
setState(c);
return free;
}

3.3 公平锁与非公平锁

公平锁与非公平锁的差别,全部集中在 tryAcquire 的一行代码上。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
// 非公平:上来就抢,不管队里有没有人
static final class NonfairSync extends Sync {
protected boolean tryAcquire(int acquires) {
return nonfairTryAcquire(acquires); // 见 3.2,直接 CAS
}
}

// 公平:先看看前面有没有排队的,有就老实排
static final class FairSync extends Sync {
protected boolean tryAcquire(int acquires) {
final Thread current = Thread.currentThread();
int c = getState();
if (c == 0) {
// hasQueuedPredecessors:队列里是否存在着比我更早的等待者
if (!hasQueuedPredecessors() &&
compareAndSetState(0, acquires)) {
setExclusiveOwnerThread(current);
return true;
}
} else if (current == getExclusiveOwnerThread()) {
int nextc = c + acquires;
if (nextc < 0) throw new Error("Maximum lock count exceeded");
setState(nextc);
return true;
}
return false;
}
}

// 判断"我前面是否有人在排队":队空 || 队首就是我自己 || 队首线程是当前线程
public final boolean hasQueuedPredecessors() {
Node t = tail, h = head, s;
return h != t && // 队列非空
((s = h.next) == null || s.thread != Thread.currentThread());
}

一个反直觉的点:公平锁的 tryLock() 依然非公平——它直接调 sync.nonfairTryAcquire(1),哪怕对象是 new ReentrantLock(true)。因为”尝试”语义本就尽力而为,要严格公平请用 tryLock(0, TimeUnit.SECONDS)。

性能取舍很清楚:非公平吞吐远高于公平。因为唤醒一个 park 的线程要经历内核态/用户态切换(数千周期),公平锁坚持唤醒队首,这段空窗期 CPU 干等;非公平锁允许此刻刚到达、本就在运行态的线程插队,直接把空窗期填满。代价是队尾线程可能饥饿。

为什么默认非公平?因为保证公平的成本高于收益:非公平吞吐可高出数倍,而饥饿在真实场景极少发生(线程终会执行完)。只有明确要求”先来后到”且临界区较长时,才值得付这笔税。

3.4 Condition:一个锁,多个等待队列

synchronized 的 wait/notify 把等待队列绑定在对象监视器上,一个锁只有一个队列,只能 notifyAll 全量唤醒。AQS 的 ConditionObject 打破了这一限制:一把 ReentrantLock 可以 newCondition() 出任意多个等待队列,实现精准唤醒。

3.4.1 await 与 signal 的底层

ConditionObject 内部是一条单向链表(firstWaiter/lastWaiter),节点复用 Node,waitStatus 为 CONDITION(-2),靠 nextWaiter 串联。流程如下:

  • await():addConditionWaiter() 入条件队列 → fullyRelease(node) 完全释放锁(可重入时 state 可能 >1,必须一次清 0)→ park 挂起 → 被 signal 后由 transferAfterCancelledWait 判定迁移方式 → acquireQueued 重新排队争锁并恢复原 state。
  • signal():把条件队列首节点 transferForSignal 到同步队列(状态 CONDITION→0、CAS 入队尾、前驱改 SIGNAL),该线程随后被正常唤醒。
  • signalAll():条件队列所有节点依次搬回同步队列。

transferAfterCancelledWait 负责区分”等待期间被中断”的两种情形:

1
2
3
4
5
6
7
8
9
10
11
final boolean transferAfterCancelledWait(Node node) {
// 情况一:正常被 signal(状态已被改成 0),CAS 回 CONDITION 失败,说明有人动了它
if (compareAndSetWaitStatus(node, Node.CONDITION, 0)) {
enq(node); // 超时或被中断:自己把自己塞回同步队列
return true; // true 表示"中断发生在 signal 之前"
}
// 情况二:已经被 signal 了(状态已是 0),等它迁移完成即可
while (!isOnSyncQueue(node))
Thread.yield(); // signal 尚未入队,让步自旋等一会
return false; // false 表示"中断发生在 signal 之后"
}

这个返回值决定了 await() 抛 InterruptedException 的时机:中断早于 signal 才抛异常,中断晚于 signal 则先恢复锁、再补上中断标记,保证不丢事件。这比 Object.wait() 的语义更严谨。

3.4.2 用 Condition 实现有界缓冲区

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
import java.util.concurrent.locks.Condition;
import java.util.concurrent.locks.ReentrantLock;

/**
* 有界缓冲区:经典生产者-消费者
* 用两个 Condition 分别挂起"生产者"和"消费者",实现精准唤醒
*/
public class BoundedBuffer<T> {
private final Object[] items;
private int putPtr, takePtr, count;
private final ReentrantLock lock = new ReentrantLock();
private final Condition notFull = lock.newCondition(); // 队列不满
private final Condition notEmpty = lock.newCondition(); // 队列不空

@SuppressWarnings("unchecked")
public BoundedBuffer(int capacity) { items = new Object[capacity]; }

public void put(T x) throws InterruptedException {
lock.lock();
try {
while (count == items.length) // 必须用 while,防止虚假唤醒
notFull.await(); // 释放锁并挂起,被唤醒后重新抢锁
items[putPtr] = x;
if (++putPtr == items.length) putPtr = 0;
count++;
notEmpty.signal(); // 只唤醒一个消费者,不用 signalAll
} finally {
lock.unlock();
}
}

@SuppressWarnings("unchecked")
public T take() throws InterruptedException {
lock.lock();
try {
while (count == 0)
notEmpty.await();
T x = (T) items[takePtr];
if (++takePtr == items.length) takePtr = 0;
count--;
notFull.signal(); // 只唤醒一个生产者
return x;
} finally {
lock.unlock();
}
}

public static void main(String[] args) {
BoundedBuffer<Integer> buf = new BoundedBuffer<>(5);
// 生产者:每 300ms 放一个数
new Thread(() -> {
for (int i = 1; i <= 10; i++) {
try { buf.put(i); System.out.println("put " + i); Thread.sleep(300); }
catch (InterruptedException e) { Thread.currentThread().interrupt(); }
}
}, "producer").start();
// 消费者:每 800ms 取一个,必然触发 notFull.await()
new Thread(() -> {
for (int i = 1; i <= 10; i++) {
try { System.out.println("take " + buf.take()); Thread.sleep(800); }
catch (InterruptedException e) { Thread.currentThread().interrupt(); }
}
}, "consumer").start();
}
}
对比项 Object.wait/notify Condition.await/signal
依赖 必须先 synchronized 拿到监视器 必须先 lock() 拿到 Lock
等待队列数 每个对象仅 1 个 每把锁可有多个 Condition
唤醒粒度 notify 随机一个 / notifyAll 全部 signal 指定队列的首个 / signalAll 该队列全部
中断语义 抛异常,区分不了中断与 signal 先后 transferAfterCancelledWait 精确区分
超时等待 支持 wait(timeout) 支持 awaitNanos / awaitUntil(绝对时间)
释放方式 释放一次监视器 fullyRelease 释放全部重入次数

3.5 ReentrantReadWriteLock 与锁降级

读写锁把”读-读”从互斥中解放出来:读读共享、读写互斥、写写互斥,适合缓存、配置中心这类读远多于写的场景。

锁降级指”持有写锁 → 获取读锁 → 释放写锁“,最终降级为读锁,保证刚写完的数据立刻能被自己读到,且期间不会被其它写线程插入。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
import java.util.concurrent.locks.ReentrantReadWriteLock;

/** 缓存示例:演示锁降级(写锁 → 读锁 → 释放写锁) */
public class CacheDemo {
private final ReentrantReadWriteLock rwl = new ReentrantReadWriteLock();
private final ReentrantReadWriteLock.ReadLock readLock = rwl.readLock();
private final ReentrantReadWriteLock.WriteLock writeLock = rwl.writeLock();
private volatile boolean cacheValid;
private Object data;

public Object getData() {
readLock.lock();
if (!cacheValid) {
// 必须先释放读锁,否则直接获取写锁会死锁(读锁不支持升级为写锁)
readLock.unlock();
writeLock.lock();
try {
if (!cacheValid) { // 双重检查,防止重复加载
data = loadFromDb();
cacheValid = true;
}
// ↓↓↓ 锁降级开始:在持有写锁的情况下获取读锁 ↓↓↓
readLock.lock();
} finally {
writeLock.unlock(); // 降级为读锁,其它读线程可进来
}
}
try {
return data; // 仍在读锁保护下返回,期间不会有写线程改它
} finally {
readLock.unlock();
}
}

private Object loadFromDb() {
System.out.println(Thread.currentThread().getName() + " 加载数据库...");
return new Object();
}
}

为什么不支持锁升级(先读后写)?读锁可能被多线程持有,若允许其中一个升级成写锁,其它读线程还浑然不觉地持有读锁,互斥语义即被破坏;且两个读线程同时升级会互相等待而死锁。所以要先释放再获取,或一开始就拿写锁。

3.6 StampedLock:乐观读

StampedLock(JDK 8)用 long stamp 作票据,提供三种模式:

模式 方法 语义 是否阻塞
写锁 writeLock() / unlockWrite(stamp) 独占 是
悲观读锁 readLock() / unlockRead(stamp) 共享 是
乐观读 tryOptimisticRead() / validate(stamp) 不加锁,读完校验是否被写过 否

乐观读的思路类似数据库的乐观锁版本号:读前取 stamp,读完用 validate(stamp) 检查期间有无写操作,没有就直接用,有则升级为悲观读锁重读一次。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
import java.util.concurrent.locks.StampedLock;

/** StampedLock 乐观读:读多写少场景下几乎无锁开销 */
public class Point {
private double x, y;
private final StampedLock sl = new StampedLock();

// 写:独占,返回 stamp 用于释放
public void move(double deltaX, double deltaY) {
long stamp = sl.writeLock();
try {
x += deltaX;
y += deltaY;
} finally {
sl.unlockWrite(stamp);
}
}

// 乐观读:不加锁,靠 validate 兜底
public double distanceFromOrigin() {
long stamp = sl.tryOptimisticRead(); // 拿到"版本戳",开销极小
double currentX = x, currentY = y; // 读字段(注意:这里要读 volatile 或用局部变量拷贝)
if (!sl.validate(stamp)) { // 期间被写过,数据可能不一致
stamp = sl.readLock(); // 升级为悲观读锁
try {
currentX = x;
currentY = y;
} finally {
sl.unlockRead(stamp);
}
}
return Math.sqrt(currentX * currentX + currentY * currentY);
}

/** 悲观读 + 条件等待的写法(StampedLock 不支持 Condition!) */
public double pessimisticDistance() {
long stamp = sl.readLock();
try {
return Math.sqrt(x * x + y * y);
} finally {
sl.unlockRead(stamp);
}
}
}

StampedLock 有两个坑:不可重入(同线程重复获取会自锁)、不支持 Condition(需要条件等待请退回 ReentrantReadWriteLock);且它的悲观读并非 AQS 实现,不能当普通 Lock 传给需要 Lock 接口的方法。

3.7 死锁排查:jstack

ReentrantLock 死锁不会像 synchronized 那样被 jstack 自动标注 “Found one Java-level deadlock”,但会打印 WAITING (parking) 与 “parking to wait for <0x...> (a ReentrantLock$NonfairSync)”,需人工串联。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
import java.util.concurrent.locks.ReentrantLock;
import java.util.concurrent.TimeUnit;

/** 复现 ReentrantLock 死锁:两个线程互相持有对方需要的锁 */
public class LockDeadlockDemo {
private static final ReentrantLock lockA = new ReentrantLock();
private static final ReentrantLock lockB = new ReentrantLock();

public static void main(String[] args) throws InterruptedException {
// 线程 1:先 A 后 B
new Thread(() -> {
lockA.lock();
try {
TimeUnit.SECONDS.sleep(1); // 制造时间窗口,让对方也拿到 B
System.out.println("T1 尝试获取 lockB...");
lockB.lock();
try { System.out.println("T1 拿到两把锁"); }
finally { lockB.unlock(); }
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
} finally { lockA.unlock(); }
}, "Thread-A-B").start();

// 线程 2:先 B 后 A
new Thread(() -> {
lockB.lock();
try {
TimeUnit.SECONDS.sleep(1);
System.out.println("T2 尝试获取 lockA...");
lockA.lock(); // 死在这里
try { System.out.println("T2 拿到两把锁"); }
finally { lockA.unlock(); }
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
} finally { lockB.unlock(); }
}, "Thread-B-A").start();

TimeUnit.SECONDS.sleep(3);
System.out.println("主线程退出,请用 jstack <pid> 观察两个 WAITING (parking) 线程");
}
}

排查步骤:jps -l 取进程号 → jstack -l <pid> > dead.txt → 搜 WAITING (parking) 记下 parking to wait for <0x...> → 再搜该地址出现在哪个线程的 locked <0x...> 之后,即可画出”谁持有了谁想要的锁”的环。修复手段通常是固定加锁顺序(按对象 hashCode 排序加锁),或用 tryLock(timeout) + 退避重试。

四、原子类家族:无锁化的第一选择

4.1 CAS 原理

CAS(Compare And Swap)是一条 CPU 原子指令:给定内存位置 V、期望旧值 A、新值 B,当且仅当 V 等于 A 时才把 V 改成 B,否则什么都不做,并返回是否成功。x86 上对应 cmpxchg,配合 lock 前缀锁总线/缓存行保证原子性。Java 层通过 sun.misc.Unsafe 暴露:

1
2
3
4
// Unsafe 中的三个核心 CAS 方法(native,最终由 CPU 指令实现)
public final native boolean compareAndSwapObject(Object o, long offset, Object expected, Object x);
public final native boolean compareAndSwapInt (Object o, long offset, int expected, int x);
public final native boolean compareAndSwapLong(Object o, long offset, long expected, long x);

JDK 9 之后 Unsafe 被逐步收口,官方替代品是 VarHandle(JDK 9)与 MemorySegment(JDK 22 的 FFM API),AtomicInteger 内部也换成了 VarHandle,语义完全一致。

4.2 CAS 的三大问题与解法

问题 描述 解决方案
ABA 值从 A 变成 B 又变回 A,CAS 检查时以为没变过 AtomicStampedReference(版本号)、AtomicMarkableReference(布尔标记)
循环开销大 竞争激烈时 CAS 长期失败,自旋空耗 CPU LongAdder 分散热点;或退化为锁;JVM 支持 pause 指令降低功耗
只能保证单变量原子性 多个变量需要一起原子更新时无能为力 封装成对象用 AtomicReference;或直接用锁

ABA 的经典场景:栈顶元素从 A 弹出换成 B 又被压回 A,某线程 CAS 判断”栈顶还是 A”而成功,但整个栈的结构其实已经变了。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
import java.util.concurrent.atomic.AtomicStampedReference;

/** ABA 问题与 AtomicStampedReference 的解决 */
public class ABADemo {
public static void main(String[] args) throws InterruptedException {
// 普通引用:无法察觉 A -> B -> A 的变化
AtomicStampedReference<Integer> ref =
new AtomicStampedReference<>(100, 1); // 初始值 100,版本戳 1

// 线程 t1:想把 100 改成 101,但它中间睡了 1 秒
Thread t1 = new Thread(() -> {
int stamp = ref.getStamp(); // 读到版本 1
System.out.println("t1 读到 stamp=" + stamp);
try { Thread.sleep(1000); } catch (InterruptedException e) { return; }
// 期望:值仍是 100 且版本仍是 1 → 失败,因为版本已被 t2 改过
boolean ok = ref.compareAndSet(100, 101, stamp, stamp + 1);
System.out.println("t1 CAS 结果 = " + ok + ",当前值 = " + ref.getReference()
+ ",当前版本 = " + ref.getStamp());
}, "t1");

// 线程 t2:制造 ABA —— 100 -> 101 -> 100,版本号从 1 变到 3
Thread t2 = new Thread(() -> {
ref.compareAndSet(100, 101, ref.getStamp(), ref.getStamp() + 1);
System.out.println("t2 第一次修改,值=" + ref.getReference() + " 版本=" + ref.getStamp());
ref.compareAndSet(101, 100, ref.getStamp(), ref.getStamp() + 1);
System.out.println("t2 第二次修改(改回 100),值=" + ref.getReference()
+ " 版本=" + ref.getStamp());
}, "t2");

t1.start();
Thread.sleep(100); // 确保 t1 先读到 stamp
t2.start();
t1.join(); t2.join();
// 输出:t1 CAS 结果 = false —— 值虽然回到 100,但版本戳不匹配,操作被拒绝
}
}

4.3 原子类家族全景

分类 类 用途
基本类型 AtomicInteger、AtomicLong、AtomicBoolean 单变量原子读写与运算
数组 AtomicIntegerArray、AtomicLongArray、AtomicReferenceArray 数组元素的原子更新(复制了数组,不影响原数组)
引用类型 AtomicReference、AtomicStampedReference、AtomicMarkableReference 对象引用的原子更新,后两者解决 ABA
字段更新器 AtomicIntegerFieldUpdater、AtomicLongFieldUpdater、AtomicReferenceFieldUpdater 反射式更新对象的 volatile 字段,省去包装对象
累加器(JDK 8) LongAdder、DoubleAdder、LongAccumulator、DoubleAccumulator 高并发求和/自定义聚合,吞吐远高于原子类
其它 AtomicLongFieldUpdater 等 见上;Striped64 是 LongAdder 的父类

4.4 AtomicInteger 的 incrementAndGet

1
2
3
4
5
6
7
8
9
10
11
12
13
// JDK 8 AtomicInteger 源码(去掉注释)
public final int incrementAndGet() {
return unsafe.getAndAddInt(this, valueOffset, 1) + 1;
}

// Unsafe.getAndAddInt:经典的 CAS 自旋循环
public final int getAndAddInt(Object o, long offset, int delta) {
int v;
do {
v = this.getIntVolatile(o, offset); // 1. volatile 读当前值
} while (!this.compareAndSwapInt(o, offset, v, v + delta)); // 2. CAS 失败就重来
return v; // 3. 返回旧值
}

JDK 9+ 语义相同,只是把 Unsafe 换成 VarHandle.compareAndSet,并在失败循环里加 Thread.onSpinWait() 提示 CPU(x86 上是 pause 指令),降低自旋功耗。

4.5 LongAdder:分段累加与伪共享

AtomicLong 高并发累加时有致命弱点:所有线程 CAS 同一个变量,只有一个能成功,其余全部自旋重试,竞争度随线程数线性恶化。

LongAdder 的思路是”分而治之”,继承自 Striped64:

  • base:竞争不激烈时直接 CAS 这个基础值,和 AtomicLong 一样快;
  • Cell[]:竞争激烈时每线程按自己的 probe 哈希映射到某个 Cell 槽位各加各的,冲突就 rehash 换槽或扩容数组(上限为 CPU 核数);
  • sum():base + 所有 Cell 求和。这不是原子快照,并发更新时可能漏掉增量,所以 LongAdder 只适合 QPS、计数这类统计场景,不适合做需要精确一致性的状态判断。
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
// Striped64 中的 Cell 定义:注意 @Contended
@sun.misc.Contended static final class Cell {
volatile long value;
Cell(long x) { value = x; }
final boolean cas(long cmp, long val) {
return UNSAFE.compareAndSwapLong(this, valueOffset, cmp, val);
}
}

// LongAdder.add 的核心逻辑(伪代码化)
public void add(long x) {
Cell[] cs; long b, v; int m; Cell c;
if ((cs = cells) != null || !casBase(b = base, b + x)) { // 先试 base
boolean uncontended = true;
if (cs == null || (m = cs.length - 1) < 0 ||
(c = cs[getProbe() & m]) == null ||
!(uncontended = c.cas(v = c.value, v + x)))
longAccumulate(x, null, uncontended); // 兜底:初始化/扩容/换槽
}
}

@Contended 是点睛之笔。CPU 以 Cache Line(通常 64 字节) 加载内存,一个 Cell 仅 16 字节,相邻两个 Cell 会落在同一条缓存行:A 改 Cell[0]、B 改 Cell[1] 虽逻辑无关,MESI 协议却会互相把对方的缓存行置为无效,产生伪共享(False Sharing),性能退化到与单变量 CAS 相当。@Contended 让 JVM 插入 128 字节填充(默认只对 JDK 内部类生效,用户类需 -XX:-RestrictContended),使每个 Cell 独占一条缓存行。

实测对比(4 核机器,32 线程各累加 1000 万次,量级仅供参考):

实现方式 耗时(约) 相对吞吐 结果准确性
synchronized 方法 4.5 s 1x 精确
ReentrantLock 1.8 s 2.5x 精确
AtomicLong 1.2 s 3.7x 精确
LongAdder 0.25 s 18x 最终一致(sum() 非原子快照)

结论很直接:高并发计数用 LongAdder,需要精确值或做 CAS 判断用 AtomicLong,低并发下 AtomicLong 更快也更省内存(后者要维护 Cell 数组)。LongAccumulator 是通用版,可传入自定义二元运算(如 Long::max)与初始值。

五、并发容器:性能与语义的权衡

5.1 同步容器 vs 并发容器

同步容器(Vector、Hashtable、Collections.synchronizedXxx)在每个方法上加 synchronized 锁住整个容器,并发容器则把锁粒度降到元素级:

对比项 同步容器(Vector、Hashtable、Collections.synchronizedXxx) 并发容器(ConcurrentHashMap、CopyOnWriteArrayList 等)
实现方式 方法上加 synchronized,锁住整个容器 分段/CAS/写时复制,锁粒度极细
并发度 同一时刻仅一个线程可访问 读读、读写(多数情况)可并发
复合操作 size() 与 get() 之间需外部加锁 ConcurrentHashMap 提供 putIfAbsent 等原子复合方法
迭代器 快速失败(fail-fast),并发修改抛 ConcurrentModificationException 弱一致(fail-safe),迭代期间允许修改
性能 低,高竞争下急剧退化 高,随线程数近似线性扩展
适用场景 遗留代码、极低并发 一切新代码

5.2 ConcurrentHashMap:JDK 8 的重写

JDK 7 用 Segment 分段锁(继承 ReentrantLock,默认 16 段),并发度上限就是段数。JDK 8 彻底重写,改为 数组 + 链表/红黑树 + CAS + synchronized 锁单个桶:

  • put:桶空时 CAS 插入(无锁);桶非空则 synchronized 锁住桶的头节点,只锁一条链表,其它桶不受影响;
  • get:全程无锁,Node 的 val、next 均为 volatile,靠 volatile 读保证可见性;
  • 计数:借鉴 LongAdder,用 baseCount + CounterCell[] 分段计数,size() 为求和结果(非精确),推荐 mappingCount()(long);
  • 扩容:transfer 按 stride 把迁移任务切块,线程做完自己那段若还有未迁移区间会继续领取,多线程可协助迁移(helpTransfer),靠 ForwardingNode(hash = MOVED = -1)标记该桶已搬走,扩容期间读写仍可正常进行;
  • 树化:链表长度 ≥ 8 且数组长度 ≥ 64 时转红黑树,≤ 6 时退化回链表。

为什么不允许 null 键值?为了消除歧义:get(key) 返回 null 时无法区分”key 不存在”与”key 存在但值为 null”。单线程下还能用 containsKey 二次确认,并发下两次调用之间状态已变,判断永远不可靠。所以 Doug Lea 直接在设计上禁止 null,让 null 唯一表示”不存在”。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.LongAdder;

/** ConcurrentHashMap 常用原子复合操作与统计 */
public class ConcurrentHashMapDemo {
public static void main(String[] args) {
ConcurrentHashMap<String, LongAdder> wordCount = new ConcurrentHashMap<>();

// 1. putIfAbsent:不存在才放,原子操作,避免"先 get 判断再 put"的竞态
wordCount.putIfAbsent("java", new LongAdder());
wordCount.get("java").increment();

// 2. computeIfAbsent:更推荐的写法,一行搞定(注意不要在 lambda 里做耗时操作,
// 因为 compute 会锁住当前桶)
wordCount.computeIfAbsent("juc", k -> new LongAdder()).increment();
wordCount.computeIfAbsent("juc", k -> new LongAdder()).increment();

// 3. 遍历:弱一致迭代器,不会抛 ConcurrentModificationException
wordCount.forEach((k, v) -> System.out.println(k + " -> " + v.sum()));

// 4. 统计:mappingCount 返回 long,比 size() 的 int 更安全
System.out.println("总 key 数 = " + wordCount.mappingCount());

// 5. 批量:reduce / search / forEach 支持并行阈值
System.out.println("总次数 = " + wordCount.reduceValuesToLong(1L, LongAdder::sum));

// 6. 注意:value 不能为 null,会直接 NPE
// wordCount.put("null-value-demo", null); // java.lang.NullPointerException
}
}

5.3 CopyOnWriteArrayList

写时复制:每次修改(add/set/remove)都复制一份新数组,改完后用 volatile 引用切换过去,读操作完全不加锁。

代价显而易见:写期间新旧两个数组同时在堆里,元素多时易触发 GC;复制是 O(n),写多即灾难;迭代器持有创建时的快照,遍历期间其它线程的修改完全不可见,即”弱一致性”。

适用场景很窄:读极多、写极少且能容忍短暂不一致,典型是监听器列表、路由表、黑白名单配置。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
import java.util.concurrent.CopyOnWriteArrayList;

/** CopyOnWriteArrayList:读多写极少场景;演示迭代器的弱一致性 */
public class CopyOnWriteDemo {
public static void main(String[] args) throws InterruptedException {
CopyOnWriteArrayList<String> listeners = new CopyOnWriteArrayList<>();
listeners.add("listener-1");
listeners.add("listener-2");

// 迭代过程中被修改:不会抛 ConcurrentModificationException
Thread writer = new Thread(() -> {
try { Thread.sleep(10); } catch (InterruptedException e) { return; }
listeners.add("listener-3"); // 复制新数组,老迭代器看不到
System.out.println("写线程已添加 listener-3");
});
writer.start();

for (String s : listeners) { // 快照迭代,输出 1、2(看不到 3)
System.out.println("遍历到:" + s);
Thread.sleep(50);
}
writer.join();
System.out.println("遍历结束后 size = " + listeners.size()); // 3
}
}

5.4 ConcurrentLinkedQueue

基于 Michael-Scott 算法的无界非阻塞队列:入队 CAS 尾节点的 next,出队 CAS head 并帮助推进,全程无锁,适合高并发且不需要阻塞语义的场景。注意 size() 需遍历,是 O(n) 操作,判空请用 isEmpty()。

5.5 BlockingQueue 家族与选型

队列 底层结构 是否有界 锁/实现 典型用途
ArrayBlockingQueue 数组 有界(构造指定) 单把 ReentrantLock + 两个 Condition 固定容量、需要背压的池化场景
LinkedBlockingQueue 链表 可选(默认 Integer.MAX_VALUE) 双锁(putLock/takeLock)分离,吞吐更高 通用任务队列,线程池默认项
PriorityBlockingQueue 堆 无界(会自动扩容) 单锁 + 自旋 CAS 扩容 需要按优先级出队的任务调度
DelayQueue PriorityQueue 无界 单锁 + Condition available 延时任务、订单超时、缓存过期
SynchronousQueue 无存储 容量为 0 CAS 双栈/双队列 直接交接,Executors.newCachedThreadPool 用
LinkedTransferQueue 链表 无界 CAS + transfer 语义 生产者需确认”已被消费者接收”
LinkedBlockingDeque 双向链表 可选 单锁 工作窃取、双端队列
DelayedWorkQueue 堆(数组) 有界/自扩容 ScheduledThreadPoolExecutor 内部专用 定时任务调度

四组插入/移除 API 的语义必须分清:

行为 抛异常 返回特殊值 阻塞 超时
插入 add(e) offer(e) put(e) offer(e, time, unit)
移除 remove() poll() take() poll(time, unit)
检查 element() peek() — —

5.5.1 用 DelayQueue 实现订单超时关闭

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
import java.util.concurrent.*;

/** 用 DelayQueue 实现订单 30 分钟未支付自动关闭 */
public class OrderTimeoutDemo {

/** 延时任务:实现 Delayed 接口 */
static class OrderTask implements Delayed {
private final String orderId;
private final long expireAt; // 到期时间戳(毫秒)

public OrderTask(String orderId, long delayMillis) {
this.orderId = orderId;
this.expireAt = System.currentTimeMillis() + delayMillis;
}

// 返回剩余延时,<=0 表示到期,可从队列取出
@Override
public long getDelay(TimeUnit unit) {
return unit.convert(expireAt - System.currentTimeMillis(), TimeUnit.MILLISECONDS);
}

// 按到期时间排序,DelayQueue 内部是优先队列,靠它决定出队顺序
@Override
public int compareTo(Delayed o) {
return Long.compare(this.expireAt, ((OrderTask) o).expireAt);
}

public String getOrderId() { return orderId; }
}

public static void main(String[] args) throws InterruptedException {
DelayQueue<OrderTask> delayQueue = new DelayQueue<>();

// 模拟下单:三个订单,分别 2s、5s、3s 后超时
delayQueue.put(new OrderTask("ORDER-001", 2_000));
delayQueue.put(new OrderTask("ORDER-002", 5_000));
delayQueue.put(new OrderTask("ORDER-003", 3_000));
System.out.println("三个订单已下单,等待超时关闭...");

// 守护线程消费到期订单
Thread closer = new Thread(() -> {
while (!Thread.currentThread().isInterrupted()) {
try {
OrderTask task = delayQueue.take(); // 只有到期才会返回
// 真实业务:先查订单状态,仍为"待支付"才关闭并释放库存
System.out.println("[" + System.currentTimeMillis() % 100_000
+ "] 订单 " + task.getOrderId() + " 超时未支付,执行关闭");
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
break;
}
}
}, "order-timeout-closer");
closer.setDaemon(true);
closer.start();

closer.join(6_000);
}
}

生产环境补充:DelayQueue 是单机内存方案,重启即丢失。更可靠的做法是延迟消息(RocketMQ/RabbitMQ 死信或 Redis ZSet 轮询)+ 数据库兜底扫描,或接入分布式调度(XXL-Job、ElasticJob)。

5.5.2 生产者消费者的三种写法

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;

/** BlockingQueue 实现生产者-消费者的三种写法对比 */
public class ProducerConsumerDemo {

/** 写法一:put/take —— 全自动阻塞,代码最少,最常用 */
static void style1() {
BlockingQueue<String> q = new ArrayBlockingQueue<>(10);
new Thread(() -> {
try {
for (int i = 0; i < 5; i++) {
q.put("item-" + i); // 满了自动等
System.out.println("生产 item-" + i);
}
} catch (InterruptedException e) { Thread.currentThread().interrupt(); }
}).start();
new Thread(() -> {
try {
for (int i = 0; i < 5; i++) {
System.out.println("消费 " + q.take()); // 空了自动等
}
} catch (InterruptedException e) { Thread.currentThread().interrupt(); }
}).start();
}

/** 写法二:offer/poll 超时 —— 可做退避、可响应关闭信号,适合需要优雅退出的服务 */
static void style2() {
BlockingQueue<String> q = new LinkedBlockingQueue<>();
volatile boolean[] running = {true};
new Thread(() -> {
while (running[0]) {
try {
String item = q.poll(1, TimeUnit.SECONDS); // 每秒醒来检查一次
if (item != null) System.out.println("消费 " + item);
} catch (InterruptedException e) { Thread.currentThread().interrupt(); break; }
}
System.out.println("消费者退出");
}).start();
try {
q.offer("a", 2, TimeUnit.SECONDS);
q.offer("b", 2, TimeUnit.SECONDS);
Thread.sleep(1200);
running[0] = false; // 发关闭信号
} catch (InterruptedException e) { Thread.currentThread().interrupt(); }
}

/** 写法三:多生产多消费 + 毒丸(Poison Pill)优雅退出 */
static void style3() throws InterruptedException {
int producers = 2, consumers = 3;
BlockingQueue<String> q = new LinkedBlockingQueue<>(20);
String POISON = "POISON_PILL";
AtomicInteger remainingConsumers = new AtomicInteger(consumers);
CountDownLatch done = new CountDownLatch(consumers);

for (int c = 0; c < consumers; c++) {
new Thread(() -> {
try {
while (true) {
String item = q.take();
if (POISON.equals(item)) {
// 收到毒丸:把自己那份"死亡信号"再放回队列,唤醒其它消费者
if (remainingConsumers.decrementAndGet() > 0) q.put(POISON);
break;
}
System.out.println(Thread.currentThread().getName() + " 消费 " + item);
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
} finally {
done.countDown();
}
}, "consumer-" + c).start();
}

for (int p = 0; p < producers; p++) {
final int pid = p;
new Thread(() -> {
for (int i = 0; i < 6; i++) {
try { q.put("P" + pid + "-item" + i); } catch (InterruptedException e) { return; }
}
}, "producer-" + p).start();
}

Thread.sleep(500);
q.put(POISON); // 主线程投递一枚毒丸即可链式关闭所有消费者
done.await();
System.out.println("所有消费者已优雅退出");
}

public static void main(String[] args) throws InterruptedException {
style1(); Thread.sleep(500);
System.out.println("--- 写法二 ---"); style2(); Thread.sleep(1500);
System.out.println("--- 写法三 ---"); style3();
}
}

六、同步协作工具类

6.1 四工具对照

工具 是否可重用 计数方向 是否阻塞等待 核心语义 典型场景
CountDownLatch 否(一次性) 递减到 0 await() 阻塞 一个或多个线程等其它线程做完 启动检查、并发压测、服务就绪
CyclicBarrier 是(可 reset) 递增到 parties await() 阻塞 一组线程互相等待,到齐一起放行 多阶段计算、并行迭代
Semaphore 是 许可加减 acquire() 阻塞 控制同时访问的线程数 限流、资源池、数据库连接数
Phaser 是(动态注册) 分阶段(phase) arriveAndAwaitAdvance 多阶段 + 参与者可动态增减 复杂流水线、替代 Barrier+Latch
Exchanger 是 成对交换 exchange() 阻塞 两个线程在汇合点交换数据 双缓冲、校对数据

6.2 CountDownLatch:一次性倒数门闩

内部是 AQS 的共享模式:state 即计数值,await() 等价于 acquireSharedInterruptibly(1),只有 state 归零才返回 1;countDown() 等价于 releaseShared(1),减到 0 时触发级联唤醒。归零后无法重置,是一次性消耗品。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;

/**
* 用 CountDownLatch 做接口并发压测
* 两个门闩:startGate 让所有线程同时发令,endGate 等所有线程跑完再统计
*/
public class ConcurrentPressureTest {
public static void main(String[] args) throws InterruptedException {
int threadCount = 100;
ExecutorService pool = Executors.newFixedThreadPool(32);
CountDownLatch startGate = new CountDownLatch(1); // 发令枪
CountDownLatch endGate = new CountDownLatch(threadCount); // 终点线
AtomicInteger success = new AtomicInteger();
AtomicInteger failure = new AtomicInteger();
AtomicLongHolder totalCost = new AtomicLongHolder();

for (int i = 0; i < threadCount; i++) {
pool.submit(() -> {
try {
startGate.await(); // 全部就位,等发令
long start = System.currentTimeMillis();
boolean ok = mockCallApi(); // 被测接口
long cost = System.currentTimeMillis() - start;
totalCost.add(cost);
if (ok) success.incrementAndGet(); else failure.incrementAndGet();
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
} finally {
endGate.countDown(); // 每跑完一个减一
}
});
}

long begin = System.currentTimeMillis();
startGate.countDown(); // 发令!
endGate.await(); // 等全部结束
long wall = System.currentTimeMillis() - begin;

System.out.println("线程数 = " + threadCount);
System.out.println("成功 = " + success.get() + ",失败 = " + failure.get());
System.out.println("总耗时 = " + wall + " ms");
System.out.println("平均 RT = " + (totalCost.sum() / threadCount) + " ms");
System.out.println("估算 QPS = " + (threadCount * 1000L / wall));
pool.shutdown();
}

/** 模拟接口调用:这里替换为真实 HTTP/RPC 调用即可 */
private static boolean mockCallApi() throws InterruptedException {
TimeUnit.MILLISECONDS.sleep(50 + (long) (Math.random() * 50));
return Math.random() > 0.05; // 95% 成功率
}

/** 简易 long 累加器(演示用,生产环境直接用 LongAdder) */
static class AtomicLongHolder {
private final LongAdder adder = new LongAdder();
void add(long v) { adder.add(v); }
long sum() { return adder.sum(); }
}
}

6.3 CyclicBarrier:可复用的栅栏

与 CountDownLatch 的两点本质区别:可循环使用(一代结束自动重置,也可 reset()),以及支持 barrierAction(所有线程到齐后、放行前,由最后一个到达的线程执行的 Runnable)。内部靠 ReentrantLock + Condition 实现,而非 AQS 共享模式。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
import java.util.concurrent.*;

/** CyclicBarrier:多线程分阶段计算,每阶段汇总一次 */
public class CyclicBarrierDemo {
public static void main(String[] args) {
int parties = 3;
int[][] matrix = {
{1, 2, 3},
{4, 5, 6},
{7, 8, 9}
};
int[] partial = new int[parties];

// barrierAction:最后到达的线程执行汇总
CyclicBarrier barrier = new CyclicBarrier(parties, () -> {
int sum = 0;
for (int p : partial) sum += p;
System.out.println("本阶段汇总结果 = " + sum);
});

ExecutorService pool = Executors.newFixedThreadPool(parties);
for (int i = 0; i < parties; i++) {
final int idx = i;
pool.submit(() -> {
try {
for (int phase = 0; phase < 2; phase++) { // 跑两轮,验证可重用
int s = 0;
for (int v : matrix[idx]) s += v + phase;
partial[idx] = s;
System.out.println(Thread.currentThread().getName()
+ " 第 " + phase + " 阶段完成,部分和=" + s);
barrier.await(); // 到齐才继续;超时可写 await(5, TimeUnit.SECONDS)
}
} catch (InterruptedException | BrokenBarrierException e) {
// BrokenBarrierException:有线程被中断/超时,栅栏被打破
System.err.println("栅栏被打破:" + e.getMessage());
}
});
}
pool.shutdown();
}
}

6.4 Semaphore:许可模型与限流器

Semaphore 同样是 AQS 共享模式,state 即剩余许可数。acquire(n) 拿 n 个许可,release(n) 归还(可多于 acquire 的数量,等于动态扩容许可)。典型用途是限流与对象池。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;

/** 基于 Semaphore 的接口限流器:QPS 上限 + 超时快速失败 + 实时统计 */
public class SemaphoreLimiter {
/** 同时允许的并发数(信号量许可数) */
private final Semaphore semaphore;
/** 是否公平:true 则先进先出,防止请求饥饿 */
private final boolean fair;

public SemaphoreLimiter(int permits, boolean fair) {
this.semaphore = new Semaphore(permits, fair);
this.fair = fair;
}

/** 阻塞式执行:拿不到许可就等 */
public <T> T execute(Callable<T> task) throws Exception {
semaphore.acquire();
try {
return task.call();
} finally {
semaphore.release(); // 必须 finally 归还,否则许可永久泄漏
}
}

/** 限时执行:拿不到就快速失败,避免雪崩 */
public <T> T execute(Callable<T> task, long timeout, TimeUnit unit) throws Exception {
if (!semaphore.tryAcquire(timeout, unit)) {
throw new RejectedExecutionException("系统繁忙,请稍后再试");
}
try {
return task.call();
} finally {
semaphore.release();
}
}

public int availablePermits() { return semaphore.availablePermits(); }
public int queueLength() { return semaphore.getQueueLength(); }

public static void main(String[] args) throws InterruptedException {
SemaphoreLimiter limiter = new SemaphoreLimiter(5, true); // 最多 5 并发
AtomicInteger ok = new AtomicInteger(), rejected = new AtomicInteger();
CountDownLatch latch = new CountDownLatch(50);
ExecutorService pool = Executors.newFixedThreadPool(50);

for (int i = 0; i < 50; i++) {
final int req = i;
pool.submit(() -> {
try {
limiter.execute(() -> {
TimeUnit.MILLISECONDS.sleep(200); // 模拟耗时业务
return "req-" + req + " done";
}, 1, TimeUnit.SECONDS); // 最多等 1 秒
ok.incrementAndGet();
} catch (RejectedExecutionException e) {
rejected.incrementAndGet(); // 被限流
} catch (Exception e) {
e.printStackTrace();
} finally {
latch.countDown();
}
});
}
latch.await();
System.out.println("成功 = " + ok.get() + ",被限流 = " + rejected.get());
System.out.println("剩余许可 = " + limiter.availablePermits()
+ ",排队线程 = " + limiter.queueLength());
pool.shutdown();
}
}

6.5 Phaser 与 Exchanger 简述

Phaser 可看作 CountDownLatch + CyclicBarrier 的增强版,支持动态注册/注销参与者(register/arriveAndDeregister)与多阶段(getPhase()),重写 onAdvance 可做阶段回调并返回 true 终止,适合分阶段流水线。

Exchanger 是双线程汇合点:exchange(V x) 阻塞直到另一线程也调用它,然后两者交换数据并返回对方的值,适用于双缓冲、数据校对。注意它两两配对,奇数个线程会有一个永远等待(除非设超时)。

七、ThreadLocal 深入

7.1 用途与基本用法

ThreadLocal 提供线程封闭:每个线程持有变量的独立副本,常用于隐式传参(用户上下文、TraceId、事务上下文)。代价是引入了一条隐藏的调用链依赖,排查时不易追踪。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;

/** ThreadLocal 基本用法:用户上下文隐式传递 */
public class ThreadLocalBasic {
private static final ThreadLocal<UserContext> CONTEXT = new ThreadLocal<>();

static class UserContext {
final String userId;
UserContext(String userId) { this.userId = userId; }
}

public static void main(String[] args) {
ExecutorService pool = Executors.newFixedThreadPool(2);
for (int i = 1; i <= 4; i++) {
final String uid = "user-" + i;
pool.submit(() -> {
try {
CONTEXT.set(new UserContext(uid)); // 入口处设置
service(); // 业务方法无需传参
} finally {
CONTEXT.remove(); // ★ 必须清理!线程池复用会串号
}
});
}
pool.shutdown();
}

static void service() { dao(); }
static void dao() {
UserContext ctx = CONTEXT.get();
System.out.println(Thread.currentThread().getName() + " 执行 SQL,uid=" + ctx.userId);
}
}

7.2 ThreadLocalMap 的结构

ThreadLocal 本身不存值,真正的容器是每个 Thread 内部的 ThreadLocalMap(Thread.threadLocals 字段),key 是 ThreadLocal 实例自身,value 是业务值。

1
2
3
4
5
6
7
8
9
10
11
static class ThreadLocalMap {
// ★ 关键:Entry 的 key 是弱引用
static class Entry extends WeakReference<ThreadLocal<?>> {
Object value; // value 是强引用!
Entry(ThreadLocal<?> k, Object v) {
super(k); // key 走 WeakReference
value = v;
}
}
private Entry[] table;
}

7.3 内存泄漏链条与”为什么 value 不是弱引用”

Entry 的 key 是弱引用:外部对 ThreadLocal 实例的强引用消失后,下次 GC 就回收 key,Entry 变成 key == null 的脏条目。但 value 仍是强引用,被 Thread → threadLocals → Entry → value 牢牢拽着。若线程不死亡(线程池核心线程几乎永不死亡),value 就永不回收——这就是完整的泄漏链条:

1
Thread(长生命周期) → ThreadLocalMap → Entry(key=null) → value(强引用,泄漏)

虽然 ThreadLocalMap 在 set/get/remove 时会顺带做启发式清理(expungeStaleEntry 清掉 key 为 null 的槽位),但这是被动的:之后若不再访问这个 Map,脏条目就一直躺着。

那为什么不把 value 也设成弱引用? 因为那样更糟:value 的唯一强引用就来自 Entry,设为弱引用一次 GC 就可能把你还在用的值清成 null,产生比泄漏更难排查的 bug。key 用弱引用是因为其生命周期由外部持有者决定(通常 static final);value 只能靠 remove() 主动断开。

7.4 线程池中的脏数据问题

比泄漏更常见的是脏数据:线程池复用线程,上次请求设的值没清理,下次请求 get() 直接拿到上一个用户的数据——严重时会造成跨用户数据泄露。三条铁律:try-finally 包裹、remove() 放 finally;入口统一设置、出口统一清理(Filter/Interceptor/AOP 最合适);不在异步子线程里直接读父线程的 ThreadLocal。

7.5 InheritableThreadLocal 与 TransmittableThreadLocal

类型 父子线程传递 线程池场景 原理
ThreadLocal 不传递 不可用 数据存在各自 Thread 的 map 里
InheritableThreadLocal 创建子线程时拷贝一次 不可用(线程池线程早已创建) Thread 构造时把父 inheritableThreadLocals 复制到子线程
TransmittableThreadLocal(TTL) 提交任务时传递 可用 阿里开源,包装 Runnable,在任务提交时刻捕获上下文,执行前注入、执行后还原

InheritableThreadLocal 只在 new Thread() 那一刻生效,而线程池的 worker 线程早已创建并被复用,父线程上下文根本传不进去。TTL 用 TtlRunnable.get(runnable) 包装任务,在 submit 时快照上下文、run 前注入、run 后还原,是目前 TraceId 异步透传的事实标准。

7.6 框架中的典型应用

  • Spring 事务:TransactionSynchronizationManager 用多个 ThreadLocal 保存当前事务的 ConnectionHolder,保证同一线程内多个 DAO 拿到同一个连接——这正是事务成立的前提;也因此事务上下文不能跨线程传播,子线程的数据库操作不参与主线程事务。
  • MDC 日志:MDC 底层是 ThreadLocal<Map<String,String>>,Filter 里 MDC.put("traceId", id)、模板加 %X{traceId} 即可全链路打印,记得 finally 里 MDC.clear()。

八、实战与面试题

8.1 AQS 与锁

Q1:AQS 的核心思想?
三件套:volatile int state 表示资源、CLH 变体的 FIFO 双向队列管排队、Node.waitStatus 表示节点状态;再用模板方法把排队/阻塞/唤醒固化在基类,只把 tryAcquire/tryRelease/tryAcquireShared/tryReleaseShared/isHeldExclusively 留给子类。

Q2:AQS 为什么用 CLH 队列变体?
一是入队必须无锁,CAS 挂尾 + 自旋重试天然适应;二是要支持节点取消,双向 prev 让 cancelAcquire 能摘除中间节点,单向做不到;三是把自旋等待改为 LockSupport.park(),避免空转烧 CPU。

Q3:SIGNAL 有什么用?为什么必须先设前驱为 SIGNAL 才能 park?
SIGNAL 表示”我释放时有义务唤醒后继”。不设就 park,前驱释放时不知道后面有人等就不会 unpark,该线程永久挂起。shouldParkAfterFailedAcquire 先 CAS 把前驱改成 SIGNAL 并返回 false 让外层再转一圈,正是为闭合这个契约。

Q4:公平锁和非公平锁的实现差异?
差异只在 tryAcquire:非公平直接 CAS 抢;公平在 CAS 前先调 hasQueuedPredecessors() 判断队中是否有更早的等待者。另注意 ReentrantLock.tryLock() 始终非公平,即便对象创建为公平锁。

Q5:为什么默认非公平锁?
唤醒 park 的线程要内核态切换,期间 CPU 空转;非公平允许此刻刚到达、本就在运行态的线程插队,填满空窗期,吞吐可高出数倍。公平锁坚持先来后到,代价是大量无效唤醒,而饥饿在真实场景极少发生。

Q6:Condition.await() 为什么必须”完全释放”锁?
锁可重入,state 可能大于 1。fullyRelease 一次把 state 清 0 并置空持有线程,否则其它线程永远拿不到锁;唤醒后再用 acquireQueued 抢锁并把 state 恢复到原重入深度。

Q7:transferAfterCancelledWait 的作用?
区分”中断发生在 signal 之前还是之后”。把 CONDITION CAS 成 0 成功,说明还没被 signal(自己超时/中断),需自己 enq 回同步队列,返回 true,await 抛 InterruptedException;失败说明已被 signal 迁移,自旋等其入队即可,返回 false,await 正常返回但补中断标记——既不丢 signal 也不吞中断。

Q8:读写锁为什么支持降级不支持升级?
降级是”写锁 → 读锁 → 释放写锁”,始终自己持有,安全。升级是”读锁 → 写锁”,而读锁可能被多线程持有,允许其中一个升级则互斥语义被破坏,且两个读线程同时升级会互相等待死锁。

Q9:synchronized 和 ReentrantLock 怎么选?
优先 synchronized(不会忘解锁、JVM 有锁消除与自适应自旋);需要可中断、超时、尝试加锁、公平性、多条件队列、读写分离时才上 ReentrantLock。

8.2 原子类与容器

Q10:CAS 的三大问题及解决方案?
ABA:用 AtomicStampedReference(版本戳)或 AtomicMarkableReference(布尔标记)。循环开销大:高并发改用 LongAdder 分散热点或退化为锁,JDK 9+ 用 Thread.onSpinWait() 降低自旋功耗。只能保证单变量原子性:把多字段封装成对象用 AtomicReference,或直接用锁。

Q11:LongAdder 为什么比 AtomicLong 快?sum() 精确吗?
base + Cell[] 分段累加,各线程哈希到不同 Cell 独立 CAS,把单点竞争拆成多槽竞争,冲突就 rehash 或扩容;Cell 上加 @Contended 填充缓存行避免伪共享。缺点:sum() 是 base 与所有 Cell 求和,并发更新时不是原子快照,只适合统计场景。低并发下 AtomicLong 更快更省内存。

Q12:ConcurrentHashMap 在 JDK 7 和 JDK 8 有什么不同?
JDK 7 用 Segment 分段锁(继承 ReentrantLock,默认 16 段,并发度被段数锁死);JDK 8 废弃 Segment,改成 Node 数组 + 链表/红黑树,空桶 CAS 插入、非空桶 synchronized 锁桶头;计数改 baseCount + CounterCell[];扩容支持多线程 helpTransfer 协助迁移,用 ForwardingNode(hash=MOVED)标记已迁移桶。

8.3 ThreadLocal

Q13:ThreadLocal 为什么会发生内存泄漏?如何避免?
Entry 的 key 是弱引用,GC 后变 null,但 value 仍是强引用,被 Thread → threadLocals → Entry → value 拽着;线程池核心线程长期存活,value 就永不回收。避免方式:用完在 finally 中 remove()。注意 expungeStaleEntry 的启发式清理是被动的,不能替代 remove()。

Q14:为什么 key 用弱引用而 value 用强引用?
key 的生命周期由外部持有者决定(通常 static final),弱引用可在持有者消失后把 Entry 标成脏条目;value 的唯一强引用就来自 Entry,若设为弱引用,GC 后可能变 null,导致业务读到 null 的诡异 bug,比泄漏更难排查。故 value 必须强引用,靠 remove() 主动断开。

Q15:线程池里用 ThreadLocal 有什么风险?
一是脏数据:线程复用,上次没清理的值被下次请求读到,可能造成跨用户数据泄露;二是内存泄漏:线程不死,value 不回收。解法是在 Filter/Interceptor/AOP 统一设置与清理,remove() 放 finally;异步子线程用 TransmittableThreadLocal 透传上下文。

8.4 编码题

Q16:三个线程按顺序打印 ABC 各 10 次(用 ReentrantLock + Condition 实现)

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
import java.util.concurrent.locks.Condition;
import java.util.concurrent.locks.ReentrantLock;

/**
* 三个线程轮流打印 A/B/C 各 10 轮
* 思路:一个 state 变量表示"该谁打印",每个线程一个 Condition 精准唤醒下一个
*/
public class PrintABC {
private static final ReentrantLock lock = new ReentrantLock();
private static final Condition cA = lock.newCondition();
private static final Condition cB = lock.newCondition();
private static final Condition cC = lock.newCondition();
private static int state = 0; // 0->打印A, 1->打印B, 2->打印C
private static final int ROUND = 10;

static class Printer implements Runnable {
private final String name;
private final int target; // 我负责的 state
private final Condition self; // 我等待的条件
private final Condition next; // 我要唤醒的条件
private final int nextState;

Printer(String name, int target, Condition self, Condition next, int nextState) {
this.name = name; this.target = target;
this.self = self; this.next = next; this.nextState = nextState;
}

@Override
public void run() {
for (int i = 0; i < ROUND; i++) {
lock.lock();
try {
while (state != target) // while 防虚假唤醒
self.await();
System.out.print(name);
state = nextState;
next.signal(); // 精准唤醒下一个,不用 signalAll
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
return;
} finally {
lock.unlock();
}
}
}
}

public static void main(String[] args) throws InterruptedException {
Thread a = new Thread(new Printer("A", 0, cA, cB, 1), "T-A");
Thread b = new Thread(new Printer("B", 1, cB, cC, 2), "T-B");
Thread c = new Thread(new Printer("C", 2, cC, cA, 0), "T-C");
a.start(); b.start(); c.start();
a.join(); b.join(); c.join();
System.out.println("\n完成"); // 输出:ABCABCABC...(共 10 轮)
}
}

Q17:实现一个支持限流与超时回退的简易资源池

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;

/** 用 Semaphore + BlockingQueue 实现一个简易数据库连接池 */
public class SimpleConnectionPool {
private final BlockingQueue<Connection> pool;
private final Semaphore permits; // 与池中连接数一致
private final AtomicInteger created = new AtomicInteger();

static class Connection { // 模拟连接
final int id;
Connection(int id) { this.id = id; }
@Override public String toString() { return "Conn#" + id; }
}

public SimpleConnectionPool(int size) {
this.pool = new ArrayBlockingQueue<>(size);
this.permits = new Semaphore(size, true); // 公平:先到先得
for (int i = 0; i < size; i++) pool.add(new Connection(i));
}

/** 借连接:最多等 timeout;拿到许可后再 take(此时队列必有货) */
public Connection borrow(long timeout, TimeUnit unit) throws InterruptedException {
if (!permits.tryAcquire(timeout, unit))
throw new IllegalStateException("连接池耗尽,等待超时");
try {
return pool.take();
} catch (InterruptedException e) {
permits.release(); // 被打断要归还许可,避免泄漏
throw e;
}
}

/** 还连接:先放回队列再释放许可 */
public void release(Connection c) {
if (c == null) return;
pool.offer(c);
permits.release();
}

public int available() { return permits.availablePermits(); }

public static void main(String[] args) throws InterruptedException {
SimpleConnectionPool p = new SimpleConnectionPool(3);
ExecutorService es = Executors.newFixedThreadPool(10);
CountDownLatch latch = new CountDownLatch(10);
for (int i = 0; i < 10; i++) {
es.submit(() -> {
Connection c = null;
try {
c = p.borrow(1, TimeUnit.SECONDS);
System.out.println(Thread.currentThread().getName() + " 借到 " + c);
TimeUnit.MILLISECONDS.sleep(300);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
} catch (IllegalStateException e) {
System.out.println(Thread.currentThread().getName() + " -> " + e.getMessage());
} finally {
if (c != null) p.release(c);
latch.countDown();
}
});
}
latch.await();
System.out.println("归还后可用连接 = " + p.available()); // 回到 3
es.shutdown();
}
}

8.5 本篇小结

JUC 的心法浓缩成三句话:能不加锁就不加锁(无锁 CAS、线程封闭、不可变对象);必须加锁就把粒度降到最小(ConcurrentHashMap 的单个桶、LongAdder 的单个 Cell,本质都是”拆热点”);锁之外的协作交给现成工具(Latch/Barrier/Semaphore/BlockingQueue)。下一篇进入 Executor 体系,看线程池如何把这些原语组装成工业级任务调度框架。