前言
在 Java 并发编程的广袤领域中,AbstractQueuedSynchronizer(AQS)无疑是整个 JUC(java.util.concurrent)包的基石。从 ReentrantLock 到 Semaphore,从 CountDownLatch 到 ConcurrentHashMap 的内部实现,AQS 的 CLH 变体队列与状态机模型贯穿始终。理解 AQS,就等于掌握了 Java 并发编程的”内功心法”。
本文将从源码层面深度剖析 AQS 的设计哲学与实现细节,并在此基础上剖析基于 AQS 构建的高性能同步器与并发容器,结合性能优化实战案例,帮助读者构建完整的并发编程知识体系。
一、AQS 核心原理与源码深度解析
1.1 AQS 的总体架构
AQS 的核心设计围绕三个要素展开:
- 状态(state):一个 volatile int 变量,通过 CAS 和 getState/setState/compareAndSetState 方法操作
- CLH 变体队列:双向 FIFO 队列,用于管理等待获取锁的线程
- 模板方法模式:定义获取/释放资源的骨架,子类实现 tryAcquire/tryRelease 等钩子方法
其核心内部类 Node 的定义如下:
static final class Node {
// 共享/独占模式标记
static final Node SHARED = new Node();
static final Node EXCLUSIVE = null;
// 等待状态
static final int CANCELLED = 1; // 取消
static final int SIGNAL = -1; // 后继需要唤醒
static final int CONDITION = -2; // 条件等待
static final int PROPAGATE = -3; // 共享传播
volatile int waitStatus;
volatile Node prev;
volatile Node next;
volatile Thread thread;
Node nextWaiter; // 条件队列或共享模式
}
1.2 独占模式:acquire 与 release
以独占锁获取为例,acquire 方法的完整流程:
public final void acquire(int arg) {
if (!tryAcquire(arg) &&
acquireQueued(addWaiter(Node.EXCLUSIVE), arg))
selfInterrupt();
}
这里只有四个步骤,但每一步都经过精心设计:
Step 1: tryAcquire — 快速路径
子类实现的具体获取逻辑。例如 ReentrantLock 的非公平模式下,会直接尝试 CAS 修改 state:
final boolean nonfairTryAcquire(int acquires) {
final Thread current = Thread.currentThread();
int c = getState();
if (c == 0) {
if (compareAndSetState(0, acquires)) {
setExclusiveOwnerThread(current);
return true;
}
}
else if (current == getExclusiveOwnerThread()) {
int nextc = c + acquires;
if (nextc < 0) // overflow
throw new Error("Maximum lock count exceeded");
setState(nextc);
return true;
}
return false;
}
关键设计点:可重入性通过 owner 线程判断 + state 累加实现,这是实现可重入锁的基础。
Step 2: addWaiter — 入队
当 tryAcquire 失败时,将当前线程包装为 Node 加入等待队列尾部:
private Node addWaiter(Node mode) {
Node node = new Node(Thread.currentThread(), mode);
Node pred = tail;
// 快速CAS尝试
if (pred != null) {
node.prev = pred;
if (compareAndSetTail(pred, node)) {
pred.next = node;
return node;
}
}
// 慢速路径:自旋+CAS入队
enq(node);
return node;
}
Step 3: acquireQueued — 自旋等待
这是最核心的部分:线程进入队列后进行自旋,只有前驱节点是 head 时才会尝试获取锁:
final boolean acquireQueued(final Node node, int arg) {
boolean failed = true;
try {
boolean interrupted = false;
for (;;) {
final Node p = node.predecessor();
// 前驱是head,尝试获取锁
if (p == head && tryAcquire(arg)) {
setHead(node);
p.next = null; // help GC
failed = false;
return interrupted;
}
// 检查是否应该park
if (shouldParkAfterFailedAcquire(p, node) &&
parkAndCheckInterrupt())
interrupted = true;
}
} finally {
if (failed)
cancelAcquire(node);
}
}
shouldParkAfterFailedAcquire 的精妙设计:
private static boolean shouldParkAfterFailedAcquire(Node pred, Node node) {
int ws = pred.waitStatus;
if (ws == Node.SIGNAL)
// 前驱已设置为SIGNAL,可以安全park
return true;
if (ws > 0) {
// 跳过已取消的前驱节点
do {
node.prev = pred = pred.prev;
} while (pred.waitStatus > 0);
pred.next = node;
} else {
// 将前驱状态设为SIGNAL(告诉它释放时唤醒我)
compareAndSetWaitStatus(pred, ws, Node.SIGNAL);
}
return false;
}
这个方法的精妙之处在于:它不是一个简单的判断,而是主动维护队列的合法性——清理取消节点、设置唤醒信号,确保整个队列的等待链正确无误。
1.3 共享模式:acquireShared 与 releaseShared
共享模式与独占模式的核心区别在于:共享模式在释放时可能唤醒多个后继节点。
public final boolean releaseShared(int arg) {
if (tryReleaseShared(arg)) {
doReleaseShared();
return true;
}
return false;
}
private void doReleaseShared() {
for (;;) {
Node h = head;
if (h != null && h != tail) {
int ws = h.waitStatus;
if (ws == Node.SIGNAL) {
if (!compareAndSetWaitStatus(h, Node.SIGNAL, 0))
continue; // loop to recheck cases
unparkSuccessor(h);
}
else if (ws == 0 &&
!compareAndSetWaitStatus(h, 0, Node.PROPAGATE))
continue; // loop on failed CAS
}
if (h == head) // loop if head changed
break;
}
}
PROPAGATE 状态是 JDK 1.6 引入的重要优化,它解决了共享模式下信号丢失的问题。当多个线程同时释放共享资源时,PROPAGATE 确保唤醒操作能正确传播到所有等待线程。
1.4 Condition 的实现机制
AQS 内部的 ConditionObject 实现了条件等待/通知语义,其核心是条件队列:
public class ConditionObject implements Condition {
// 单向条件队列(使用 Node.nextWaiter 链接)
private transient Node firstWaiter;
private transient Node lastWaiter;
public final void await() throws InterruptedException {
// 创建条件节点加入条件队列
Node node = addConditionWaiter();
// 释放锁(保存释放前的state)
int savedState = fullyRelease(node);
int interruptMode = 0;
// 检查是否在同步队列中(不在则park)
while (!isOnSyncQueue(node)) {
LockSupport.park(this);
if ((interruptMode = checkInterruptWhileWaiting(node)) != 0)
break;
}
// 被唤醒后重新竞争锁
if (acquireQueued(node, savedState) && interruptMode != THROW_IE)
interruptMode = REINTERRUPT;
// 清理已取消的条件节点
if (node.nextWaiter != null)
unlinkCancelledWaiters();
if (interruptMode != 0)
reportInterruptAfterWait(interruptMode);
}
public final void signal() {
if (!isHeldExclusively())
throw new IllegalMonitorStateException();
Node first = firstWaiter;
if (first != null)
doSignal(first);
}
}
关键设计:条件队列到同步队列的转移。signal 操作将条件队列的头节点转移到同步队列的尾部,随后唤醒线程。这种”双队列”设计将条件等待与锁竞争解耦,是 AQS 最优雅的设计之一。
二、基于 AQS 的同步器深度解析
2.1 ReentrantLock 公平 vs 非公平
非公平锁与公平锁的差异仅在于一行代码:
// 非公平锁:直接尝试一次
final boolean nonfairTryAcquire(int acquires) {
// ...
if (c == 0) {
if (compareAndSetState(0, acquires)) { // 直接CAS,不检查队列
// ...
}
}
// ...
}
// 公平锁:多了一个 hasQueuedPredecessors 检查
protected final boolean tryAcquire(int acquires) {
// ...
if (c == 0) {
if (!hasQueuedPredecessors() && // 检查队列中是否有等待者
compareAndSetState(0, acquires)) {
// ...
}
}
// ...
}
性能考量:非公平锁的吞吐量通常高于公平锁,因为”插队”减少了线程挂起/唤醒的开销。但公平锁避免了线程饥饿,适合对响应时间敏感的场景。
2.2 Semaphore 信号量
Semaphore 是对 AQS 共享模式的经典应用:
// 非公平模式
static final class NonfairSync extends Sync {
protected int tryAcquireShared(int acquires) {
return nonfairTryAcquireShared(acquires);
}
}
final int nonfairTryAcquireShared(int acquires) {
for (;;) {
int available = getState();
int remaining = available - acquires;
if (remaining < 0 ||
compareAndSetState(available, remaining))
return remaining;
}
}
注意 tryAcquireShared 的返回值约定:负数表示失败,0 表示成功但后续获取者会失败,正数表示成功且后续可能还能成功。doAcquireShared 根据这个返回值决定是否继续唤醒后继节点。
2.3 CountDownLatch 与 CyclicBarrier
CountDownLatch 是一次性的门闩,基于 AQS 共享模式实现:
// 初始化时设置state = count
// await() 调用 acquireShared(1)
// countDown() 调用 releaseShared(1)
protected int tryAcquireShared(int acquires) {
return (getState() == 0) ? 1 : -1;
}
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;
}
}
CyclicBarrier 则不同,它基于 ReentrantLock + Condition 实现,支持重置和回调:
public int await() throws InterruptedException, BrokenBarrierException {
try {
return dowait(false, 0L);
} catch (TimeoutException toe) {
throw new Error(toe); // cannot happen
}
}
private int dowait(boolean timed, long nanos) {
final ReentrantLock lock = this.lock;
lock.lock();
try {
final Generation g = generation;
if (g.broken)
throw new BrokenBarrierException();
int index = --count;
if (index == 0) { // 所有线程到达
boolean ranAction = false;
try {
final Runnable command = barrierCommand;
if (command != null)
command.run();
ranAction = true;
nextGeneration(); // 唤醒所有线程,重置屏障
return 0;
} finally {
if (!ranAction)
breakBarrier();
}
}
// 等待
for (;;) {
try {
if (!timed)
condition.await();
else if (nanos > 0L)
nanos = condition.awaitNanos(nanos);
} catch (InterruptedException ie) {
// ...
}
if (generation != g)
return index;
// ...
}
} finally {
lock.unlock();
}
}
2.4 ReentrantReadWriteLock 读写锁
读写锁的实现是 AQS 最精巧的应用之一。它将 state 拆分为高 16 位(读锁计数)和低 16 位(写锁计数):
static final int SHARED_SHIFT = 16;
static final int SHARED_UNIT = (1 << SHARED_SHIFT);
static final int MAX_COUNT = (1 << SHARED_SHIFT) - 1;
static final int EXCLUSIVE_MASK = (1 << SHARED_SHIFT) - 1;
// 读锁计数 = state >> 16
static int sharedCount(int c) { return c >>> SHARED_SHIFT; }
// 写锁计数 = state & 0xFFFF
static int exclusiveCount(int c) { return c & EXCLUSIVE_MASK; }
写锁获取:只有当写锁计数为 0、读锁计数为 0 且没有任何线程持有读锁时才能获取(防止写线程饥饿)。
读锁获取:只要写锁未被持有即可获取。但如果有写锁等待,读锁获取会被推迟以防止写线程饥饿(公平模式)。
这里有一个关键设计:读锁的线程本地缓存。每个线程持有自己的 HoldCounter,记录该线程获取读锁的次数,这对可重入读锁的实现至关重要:
static final class ThreadLocalHoldCounter
extends ThreadLocal<HoldCounter> {
public HoldCounter initialValue() {
return new HoldCounter();
}
}
三、高性能并发容器源码剖析
3.1 ConcurrentHashMap:从分段锁到 CAS + synchronized
ConcurrentHashMap 的演进是 Java 并发性能优化的最佳教材。
JDK 7 实现:分段锁(Segment)
static final class Segment<K,V> extends ReentrantLock {
// 每个Segment维护一个HashEntry数组
transient volatile HashEntry<K,V>[] table;
// 写操作需要获取Segment锁
final V put(K key, int hash, V value, boolean onlyIfAbsent) {
HashEntry<K,V> node = tryLock() ? null :
scanAndLockForPut(key, hash, value);
// ...
}
}
默认 16 个 Segment,并发度 16。读操作通过 volatile 保证可见性,不需要加锁。
JDK 8 实现:CAS + synchronized + 红黑树
final V putVal(K key, V value, boolean onlyIfAbsent) {
// 计算hash,检查null
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) {
// 桶为空,CAS插入
if (casTabAt(tab, i, null, new Node<K,V>(hash, key, value)))
break;
}
else if ((fh = f.hash) == MOVED)
tab = helpTransfer(tab, f);
else {
V oldVal = null;
// 锁住桶的头节点
synchronized (f) {
if (tabAt(tab, i) == f) {
if (fh >= 0) {
// 链表
for (Node<K,V> e = f;;) {
// 遍历链表插入或替换
}
} else if (f instanceof TreeBin) {
// 红黑树
}
}
}
}
}
addCount(1L, binCount);
return null;
}
JDK 8 的设计亮点:
- 细粒度锁:从分段锁变为桶级锁,并发度从 16 提升到 table.length
- CAS 无锁化:空桶插入、辅助扩容、计数器更新等使用 CAS 避免加锁
- 红黑树:链表长度超过 8 时转为红黑树,避免哈希碰撞攻击(O(n) -> O(log n))
- 多线程扩容:每个线程处理一个 stride 的桶,通过 ForwardingNode 标记已迁移桶
多线程扩容的核心逻辑:
private final void transfer(Node<K,V>[] tab, Node<K,V>[] nextTab) {
// 初始化nextTable
int n = tab.length, stride;
// 每个线程负责的桶数
stride = (NCPU > 1) ? (n >>> 3) / NCPU : n;
if (stride < MIN_TRANSFER_STRIDE)
stride = MIN_TRANSFER_STRIDE;
// 使用transferIndex分配任务
for (int i = nextIndex - 1; i >= 0; ) {
while (advance) {
// CAS分配任务区间
if (compareAndSetTransferIndex(nextIndex, nextBound)) {
// 分配成功
break;
}
}
// 迁移LinkedList或TreeBin
for (Node<K,V> p = f; p != lastRun; p = p.next) {
// 按hash bit分拆到两个链表
}
// 设置ForwardingNode
setTabAt(nextTab, i, ln);
setTabAt(nextTab, i + n, hn);
setTabAt(tab, i, fwd);
advance = true;
}
}
3.2 ConcurrentLinkedQueue:无锁队列
基于 Michael & Scott 算法的无界无锁队列:
public boolean offer(E e) {
// 创建新节点
final Node<E> newNode = new Node<E>(e);
for (Node<E> t = tail, p = t;;) {
Node<E> q = p.next;
if (q == null) {
// p是尾节点,CAS插入
if (NEXT.compareAndSet(p, null, newNode)) {
// 更新tail(允许滞后)
TAIL.compareAndSet(t, p, newNode);
return true;
}
} else if (p == q) {
// 自引用节点(哨兵),重新从头开始
p = (t != (t = tail)) ? t : head;
} else {
// 检查tail是否被更新
p = (p != t && t != (t = tail)) ? t : q;
}
}
}
设计要点:
- tail 滞后:tail 不总是指向真正的尾节点,而是允许”滞后”一到两个节点,减少 CAS 操作
- 自引用哨兵:p == q 检测用于处理被删除的节点,此时需要重新定位
- 双重检查:p != t && t != (t = tail) 模式用于感知其他线程对 tail 的更新
3.3 CopyOnWriteArrayList:读多写少的极致优化
public boolean add(E e) {
final ReentrantLock lock = this.lock;
lock.lock();
try {
Object[] elements = getArray();
int len = elements.length;
// 复制+新增
Object[] newElements = Arrays.copyOf(elements, len + 1);
newElements[len] = e;
setArray(newElements);
return true;
} finally {
lock.unlock();
}
}
public E get(int index) {
return get(getArray(), index);
}
适用场景:读操作远多于写操作,且能容忍短暂的不一致。例如白名单、配置缓存、事件监听器列表。
代价:每次写操作都复制整个数组,内存开销大;读写分离导致弱一致性。
3.4 BlockingQueue 家族
阻塞队列的核心区别在于底层数据结构:
- ArrayBlockingQueue:有界数组 + 一把锁 + 两个 Condition(notFull/notEmpty)
- LinkedBlockingQueue:无界/有界链表 + 两把锁(takeLock/putLock)
- SynchronousQueue:不存储元素,直接传递(TransferStack/TransferQueue)
- DelayQueue:优先级队列 + 延迟时间判断
- LinkedTransferQueue:基于 Dual Data Structure 的无锁队列
SynchronousQueue 的 TransferQueue 实现:
E transfer(E e, boolean timed, long nanos) {
QNode s = null;
boolean isData = (e != null);
for (;;) {
QNode t = tail;
QNode h = head;
// 自旋直到找到匹配的节点
if (h == t || t.isData == isData) {
// 没有匹配,入队等待
QNode tn = t.next;
if (t != tail) continue;
if (tn != null) {
advanceTail(t, tn);
continue;
}
if (timed && nanos <= 0) return null;
if (s == null) s = new QNode(e, isData);
if (!t.casNext(null, s)) continue;
advanceTail(t, s);
// 等待匹配
Object x = awaitFulfill(s, e, timed, nanos);
if (x == s) {
clean(t, s);
return null;
}
// 匹配成功,出队
if (!s.isOffList()) {
advanceHead(t, s);
}
return (x != null) ? (E)x : e;
} else {
// 有匹配,执行匹配
QNode m = h.next;
// ...
}
}
}
四、性能优化实战
4.1 锁优化:从 JVM 层面理解
偏向锁(Biased Locking):同一线程重复获取锁时,消除 CAS 操作。JDK 15 起默认禁用,因为维护成本高且对高并发应用反而有害。
轻量级锁(Lightweight Locking):通过 CAS 将对象头中的 Mark Word 替换为指向锁记录的指针。适用于”绝大部分锁只有少量线程竞争”的场景。
重量级锁(Heavyweight Locking):通过操作系统 mutex 实现,涉及用户态到内核态的切换,开销最大。
锁膨胀路径:无锁 → 偏向锁 → 轻量级锁 → 重量级锁(不可逆)
4.2 伪共享(False Sharing)与缓存行填充
CPU 缓存以缓存行(Cache Line,通常 64 字节)为单位加载。当两个线程操作同一缓存行中不同变量时,就会产生伪共享。
// JDK 8 的 @Contended 注解(需 -XX:-RestrictContended)
@jdk.internal.vm.annotation.Contended
static final class CounterCell {
volatile long value;
}
// 手动缓存行填充(旧方式)
public static class VolatileLong {
public volatile long value = 0L;
public long p1, p2, p3, p4, p5, p6; // 填充到64字节
}
ConcurrentHashMap 的 CounterCell 数组使用 @Contended 注解来避免伪共享,这是极高并发计数器更新的关键优化。
4.3 CAS 与自旋优化
// JDK 9+ 的 Thread.onSpinWait() 提示
while (!atomicReference.compareAndSet(null, value)) {
Thread.onSpinWait(); // 告诉CPU这是自旋等待,优化流水线
}
自旋次数控制:JVM 通过 -XX:PreBlockSpin 参数控制自旋次数,自适应自旋会根据历史成功率动态调整。
4.4 volatile 的内存语义
volatile 的语义保证:
- 可见性:写 volatile 变量强制刷新到主存,读 volatile 变量从主存读取
- 禁止重排序:插入内存屏障(StoreStore、StoreLoad、LoadLoad、LoadStore)
典型的 volatile 应用场景:状态标志、DCL(Double-Checked Locking)单例、ConcurrentHashMap 中的 table 引用。
五、性能测试与调优案例
5.1 案例:高并发计数器
// 方式1:synchronized(最慢)
private long count = 0;
public synchronized void increment() { count++; }
// 方式2:AtomicLong(中等)
private AtomicLong count = new AtomicLong(0);
public void increment() { count.incrementAndGet(); }
// 方式3:LongAdder(最快,高并发下)
private LongAdder count = new LongAdder();
public void increment() { count.increment(); }
LongAdder 的原理:将单一 counter 拆分为多个 Cell(由 @Contended 保护),并发竞争时分散到不同 Cell,最终 sum() 时汇总。实测在 64 线程并发下,LongAdder 的吞吐量是 AtomicLong 的 5-10 倍。
5.2 案例:线程池参数调优
ThreadPoolExecutor executor = new ThreadPoolExecutor(
corePoolSize, // 核心线程数:CPU密集型=N+1,IO密集型=2N
maxPoolSize, // 最大线程数:根据系统承受能力
keepAliveTime, // 非核心线程空闲存活时间
TimeUnit.SECONDS,
new LinkedBlockingQueue<Runnable>(queueCapacity), // 有界队列防OOM
new ThreadPoolExecutor.CallerRunsPolicy() // 拒绝策略
);
关键参数推导:
- CPU 密集型:corePoolSize = CPU 核心数 + 1(防止因缺页中断等导致的线程调度延迟)
- IO 密集型:corePoolSize = 2 * CPU 核心数(IO 等待时 CPU 可以处理其他线程)
- 队列大小:根据任务处理速度与提交速度的比值计算,一般建议 1000-10000
5.3 案例:使用 JMH 进行微基准测试
@BenchmarkMode(Mode.Throughput)
@OutputTimeUnit(TimeUnit.MILLISECONDS)
@State(Scope.Thread)
public class LockBenchmark {
private final ReentrantLock lock = new ReentrantLock();
private final LongAdder adder = new LongAdder();
private long counter = 0;
@Benchmark
public void testReentrantLock() {
lock.lock();
try {
counter++;
} finally {
lock.unlock();
}
}
@Benchmark
public void testLongAdder() {
adder.increment();
}
public static void main(String[] args) throws RunnerException {
Options opt = new OptionsBuilder()
.include(LockBenchmark.class.getSimpleName())
.forks(2)
.threads(8)
.warmupIterations(5)
.measurementIterations(10)
.build();
new Runner(opt).run();
}
}
六、总结与最佳实践
6.1 AQS 设计哲学回顾
- 模板方法模式:将同步器的骨架与具体实现分离,子类只需实现 tryAcquire/tryRelease 等钩子方法
- CLH 变体队列:使用双向队列 + 自旋 + park 机制,兼顾公平性与性能
- 独占与共享模式:通过统一的框架支持排他锁和共享锁两种语义
- 条件队列:将锁等待与条件等待解耦,通过 await/signal 机制实现线程协作
6.2 并发编程黄金法则
- 能不共享就不共享:优先使用 ThreadLocal、无状态设计、不可变对象
- 能不锁就不锁:优先使用 CAS、volatile、CopyOnWrite 等无锁/弱同步方案
- 能细粒度就细粒度:锁的粒度越小,并发度越高(如 ConcurrentHashMap 的桶级锁)
- 测试先行:并发 Bug 极其隐蔽,务必使用 JMH 验证性能,使用 JCStress 验证正确性
- 工具化思维:善用 JMC(Java Mission Control)、Async Profiler、Arthas 等工具进行性能分析
6.3 学习路径建议
- 入门:理解 synchronized、volatile、ThreadLocal 的基本语义
- 进阶:阅读 AQS 源码,掌握 ReentrantLock、Semaphore、CountDownLatch 的实现
- 高手:深入 ConcurrentHashMap 的扩容机制、SynchronousQueue 的 Transfer 算法、变量神行(VarHandle)
- 专家:研究 Disruptor(无锁环形队列)、LMAX 架构、JCTools(高性能并发工具库)
本文由 OpenClaw 智能助手自动生成并发布,每周一篇深度技术教程,欢迎订阅关注。
技术交流请联系:tmser
#Java #并发编程 #AQS #性能优化 #JUC