返回文章列表
JUC并发编程
JUCAQSReentrantLock读写锁源码原理

13AQS与显式锁原理

上一篇我们学习了 JUC 锁与同步工具的用法,本篇进入原理层面:AQS(AbstractQueuedSynchronizer)到底是怎么工作的?ReentrantLock 的公平/非公平、可重入、可打断、条件变量在源码层面如何实现?读写锁的共享唤醒为什么要传播?Semaphore 的 PROPAGATE 状态解决了什么 bug?最后再看 CountDownLatch 如何复用同一套共享模式。本篇是读 JUC 源码的入门篇,建议对照 JDK 8 源码一起看。

AQS 概述

AQS 是什么

AQS 全称 AbstractQueuedSynchronizer,是阻塞式锁和相关同步器工具的框架,JUC 里的锁和同步工具几乎都建立在它之上:

基于 AQS 的并发工具家族

它的特点可以概括为:

  • 用 state 属性表示资源的状态,分独占模式和共享模式。子类需要自己定义如何维护这个状态,控制如何获取锁和释放锁:
    • getState:获取 state
    • setState:设置 state
    • compareAndSetState:CAS 机制设置 state
  • 独占模式只有一个线程能访问资源(如 ReentrantLock、写锁),共享模式允许多个线程同时访问(如 Semaphore、读锁、CountDownLatch)
  • 提供基于 FIFO 的等待队列,类似于 synchronized Monitor 中的 EntryList
  • 用条件变量实现等待/唤醒机制,支持多个条件变量,类似于 Monitor 中的 WaitSet

子类主要实现下面这些方法(AQS 中默认抛出 UnsupportedOperationException):

  • tryAcquire:独占方式获取
  • tryRelease:独占方式释放
  • tryAcquireShared:共享方式获取
  • tryReleaseShared:共享方式释放
  • isHeldExclusively:判断是否当前线程独占

AQS 顶层负责的是通用骨架:

// 如果获取锁失败
if (!tryAcquire(arg)) {
    // 入队, 可以选择阻塞当前线程  park unpark
}
 
// 如果释放锁成功
if (tryRelease(arg)) {
    // 让阻塞线程恢复运行
}

手写一个不可重入锁

理解一个框架最好的方式,是先用它实现一个最简单的东西。下面自定义同步器 MySync,用 state 的 0/1 表示无锁/有锁:

final class MySync extends AbstractQueuedSynchronizer {
 
    @Override
    protected boolean tryAcquire(int acquires) {
        if (acquires == 1) {
            // CAS 抢锁: state 从 0 改为 1 才成功
            if (compareAndSetState(0, 1)) {
                setExclusiveOwnerThread(Thread.currentThread());
                return true;
            }
        }
        return false;
    }
 
    @Override
    protected boolean tryRelease(int acquires) {
        if (acquires == 1) {
            if (getState() == 0) {
                throw new IllegalMonitorStateException();
            }
            setExclusiveOwnerThread(null);
            setState(0);
            return true;
        }
        return false;
    }
 
    protected Condition newCondition() {
        // 直接复用 AQS 提供的条件变量实现
        return new ConditionObject();
    }
 
    @Override
    protected boolean isHeldExclusively() {
        return getState() == 1;
    }
}

有了自定义同步器,再包装成一个功能完备的 Lock 就非常容易,获取失败入队、可打断、超时、条件变量全部由 AQS 提供:

class MyLock implements Lock {
 
    static MySync sync = new MySync();
 
    @Override
    // 尝试,不成功,进入等待队列
    public void lock() {
        sync.acquire(1);
    }
 
    @Override
    // 尝试,不成功,进入等待队列,可打断
    public void lockInterruptibly() throws InterruptedException {
        sync.acquireInterruptibly(1);
    }
 
    @Override
    // 尝试一次,不成功返回,不进入队列
    public boolean tryLock() {
        return sync.tryAcquire(1);
    }
 
    @Override
    // 尝试,不成功,进入等待队列,有时限
    public boolean tryLock(long time, TimeUnit unit) throws InterruptedException {
        return sync.tryAcquireNanos(1, unit.toNanos(time));
    }
 
    @Override
    public void unlock() {
        sync.release(1);
    }
 
    @Override
    public Condition newCondition() {
        return sync.newCondition();
    }
}

用两个线程测试,t1 先拿锁睡 1 秒,t2 在队列中等待,t1 释放后 t2 立即获得:

MyLock lock = new MyLock();
 
new Thread(() -> {
    lock.lock();
    try {
        log.debug("locking...");
        sleep(1);
    } finally {
        log.debug("unlocking...");
        lock.unlock();
    }
}, "t1").start();
 
new Thread(() -> {
    lock.lock();
    try {
        log.debug("locking...");
    } finally {
        log.debug("unlocking...");
        lock.unlock();
    }
}, "t2").start();

输出:

22:29:28.727 c.TestAqs [t1] - locking...
22:29:29.732 c.TestAqs [t1] - unlocking...
22:29:29.732 c.TestAqs [t2] - locking...
22:29:29.732 c.TestAqs [t2] - unlocking...

这把锁是不可重入的:同一个线程连续调用两次 lock(),第二次 CAS 0→1 失败,自己把自己挡住(只会打印一次 locking):

lock.lock();
log.debug("locking...");
lock.lock();        // 自己也会被阻塞
log.debug("locking...");

实现心得与设计要点

起源:早期程序员会自己用一种同步器去实现另一种相近的同步器,例如用可重入锁实现信号量,或反之。这显然不够优雅,于是在 JSR166(Java 规范提案)中创建了 AQS,提供这种通用的同步器机制。

功能目标:

  • 阻塞版本获取 acquire 和非阻塞版本 tryAcquire
  • 获取锁超时机制
  • 通过打断取消获取
  • 独占机制与共享机制
  • 条件不满足时的等待机制

性能目标不是某一次加解锁有多快,而是可伸缩性(scalability):即使在同步器存在竞争时,也能稳定地保持效率。

AQS 的基本思想其实很简单:

// 获取锁的逻辑
while (state 状态不允许获取) {
    if (队列中还没有此线程) {
        入队并阻塞;
    }
}
当前线程出队;
 
// 释放锁的逻辑
if (state 状态允许了) {
    恢复阻塞的线程(s);
}

要点有三:原子维护 state、阻塞及恢复线程、维护队列。

  1. state 设计:state 使用 volatile 配合 CAS 保证修改的原子性;用 32 bit int 维护同步状态(当年 long 在很多平台下的测试结果不理想)。
  2. 阻塞恢复设计:早期的 suspend/resume 不可用——如果先调用 resume,再 suspend,线程将永远感知不到这次唤醒。解决方案是 park/unpark:先 unpark 再 park 也没有问题(许可机制);park/unpark 针对线程而非同步器,粒度更精细;park 还可以通过 interrupt 打断。
  3. 队列设计:FIFO 先入先出队列,不支持优先级;设计上借鉴了 CLH 队列——一种单向无锁队列,入队只需要保证 tail 赋值的原子性:
// CLH 入队伪代码
do {
    Node prev = tail;       // 原来的 tail
} while (tail.compareAndSet(prev, node));  // CAS 把 tail 改为 node
 
// CLH 出队伪代码: 自旋检查前驱节点状态
while ((Node prev = node.prev).state != 唤醒状态) {
}
head = node;

CLH 的好处是无锁、自旋、快速无阻塞。AQS 在它基础上做了改进:变成双向链表(方便从尾部向前找到需要唤醒的节点、方便取消出队),并用一个 Dummy(哑元/哨兵)节点占位。入队方法 enq 是整个队列设计的核心:

private Node enq(final Node node) {
    for (;;) {
        Node t = tail;
        // 队列还没有元素, tail 为 null
        if (t == null) {
            // 先 CAS 创建 Dummy 哨兵节点作为 head
            if (compareAndSetHead(new Node()))
                tail = head;
        } else {
            // node.prev 指向原 tail
            node.prev = t;
            // CAS 把 tail 从原 tail 改为 node
            if (compareAndSetTail(t, node)) {
                // 双向链表: 原 tail.next 指向 node
                t.next = node;
                return t;
            }
        }
    }
}

ReentrantLock 原理

类结构

ReentrantLock 内部持有一个继承自 AQS 的 Sync,它有 NonfairSync(非公平)和 FairSync(公平)两个子类;条件变量直接复用 AQS 的内部类 ConditionObject;等待队列中的节点是 AQS 的 Node。

非公平锁:加锁流程

默认构造器就是非公平锁:

public ReentrantLock() {
    sync = new NonfairSync();
}

没有竞争时很简单:lock() 里 CAS 把 state 从 0 改为 1,成功后把 exclusiveOwnerThread 设为自己。

第一个竞争出现时,Thread-0 已经持有锁(state=1),Thread-1 执行:

  1. CAS 尝试把 state 由 0 改为 1,失败;
  2. 进入 tryAcquire 逻辑,这时 state 已经是 1,仍然失败;
  3. 进入 addWaiter 逻辑,构造 Node 加入队列。Node 的创建是懒惰的,第一个节点是 Dummy 哨兵,不关联线程,仅用来占位:

第一个竞争线程构造 Dummy 节点入队

接着当前线程进入 acquireQueued 逻辑:

  1. 在死循环中不断尝试获得锁,失败后 park 阻塞;
  2. 如果自己紧邻 head(排第二位),会再次 tryAcquire,但 state 仍为 1,失败;
  3. 进入 shouldParkAfterFailedAcquire,把前驱节点(即 head)的 waitStatus 改为 -1(Node.SIGNAL),第一次返回 false;
  4. 回到 acquireQueued 再试一次,仍然失败;
  5. 再次进入 shouldParkAfterFailedAcquire,发现前驱已经是 -1,这次返回 true;
  6. 进入 parkAndCheckInterrupt,Thread-1 park(图中灰色表示阻塞)。

多个线程重复上述过程后,队列变成下面的样子——每个阻塞节点的前驱 waitStatus 都是 -1,意思是"我的后继需要被 unpark":

多个竞争失败的线程在 AQS 队列中排队阻塞

注意:是否需要 unpark,是由当前节点的前驱节点的 waitStatus == Node.SIGNAL 决定的,而不是由本节点自己的 waitStatus 决定。

非公平锁:解锁与唤醒时序

Thread-0 调用 unlock 进入 tryRelease:成功后把 exclusiveOwnerThread 置为 null、state 置 0。此时队列不为 null 且 head 的 waitStatus 为 -1,进入 unparkSuccessor:

  1. 先尝试把 head 的 waitStatus 从 -1 CAS 回 0;
  2. 从 head.next 开始找队列中离 head 最近的一个没有取消的节点(如果后继已取消,则从 tail 向前找),unpark 它——本例中即 Thread-1;
  3. Thread-1 在 acquireQueued 的 parkAndCheckInterrupt() 处恢复运行,再走一轮 for 循环:自己是老二,tryAcquire 成功;
  4. 设置 exclusiveOwnerThread = Thread-1、state=1,把自己所在节点设为新 head(清空其中的 thread),原来的 Dummy head 从链表断开,可被 GC 回收:

被唤醒的 Thread-1 成为新的 head,旧哨兵节点被断开

非公平性的体现:Thread-0 释放锁、Thread-1 尚未抢到的这个瞬间,如果新来一个 Thread-4 直接 CAS 抢占,它不需要排队:

非公平锁:新来的 Thread-4 可以插队抢锁

如果被 Thread-4 占了先,state 和独占线程都归 Thread-4,Thread-1 再次循环失败,重新把前驱置为 SIGNAL 并 park,继续排队。非公平的代价是队列中的线程可能饿一会儿,收益是减少了线程切换、吞吐量更高。

加锁核心源码

把整条链路串起来(省略了与主线无关的细节):

static final class NonfairSync extends Sync {
 
    // 加锁入口
    final void lock() {
        // 先 CAS 试一次: state 0 -> 1, 成功即拿到独占锁
        if (compareAndSetState(0, 1))
            setExclusiveOwnerThread(Thread.currentThread());
        else
            acquire(1);
    }
 
    // AQS 的模板方法
    public final void acquire(int arg) {
        if (!tryAcquire(arg) &&
            // tryAcquire 失败: 先入队 addWaiter, 再在队列中 acquireQueued
            acquireQueued(addWaiter(Node.EXCLUSIVE), arg)) {
            // 阻塞期间被中断过, 在拿到锁后补一个中断
            selfInterrupt();
        }
    }
 
    protected final boolean tryAcquire(int acquires) {
        return nonfairTryAcquire(acquires);
    }
 
    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;
            }
        }
        // 持锁线程就是自己: 锁重入, state++
        else if (current == getExclusiveOwnerThread()) {
            int nextc = c + acquires;
            if (nextc < 0)
                throw new Error("Maximum lock count exceeded");
            setState(nextc);
            return true;
        }
        return false;
    }
 
    // 把当前线程包装成独占节点接到队尾, CAS 一次失败就走 enq 自旋
    private Node addWaiter(Node mode) {
        Node node = new Node(Thread.currentThread(), mode);
        Node pred = tail;
        if (pred != null) {
            node.prev = pred;
            if (compareAndSetTail(pred, node)) {
                pred.next = node;
                return node;
            }
        }
        enq(node);
        return node;
    }
 
    // 在队列中自旋抢锁, 抢不到就 park
    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;
                }
                if (shouldParkAfterFailedAcquire(p, node) &&
                    parkAndCheckInterrupt()) {
                    // park 期间被中断过, 只记录标记, 不抛异常
                    interrupted = true;
                }
            }
        } finally {
            if (failed)
                cancelAcquire(node);
        }
    }
 
    private static boolean shouldParkAfterFailedAcquire(Node pred, Node node) {
        int ws = pred.waitStatus;
        if (ws == Node.SIGNAL) {
            // 前驱已经是 SIGNAL, 自己可以安心阻塞
            return true;
        }
        if (ws > 0) {
            // 前驱已取消, 一路向前跳过所有取消节点
            do {
                node.prev = pred = pred.prev;
            } while (pred.waitStatus > 0);
            pred.next = node;
        } else {
            // 第一次进来: 把前驱 CAS 置为 SIGNAL, 先不阻塞, 回去再抢一次
            compareAndSetWaitStatus(pred, ws, Node.SIGNAL);
        }
        return false;
    }
 
    private final boolean parkAndCheckInterrupt() {
        LockSupport.park(this);
        // interrupted() 会清除打断标记
        return Thread.interrupted();
    }
}

解锁核心源码

public final boolean release(int arg) {
    if (tryRelease(arg)) {
        Node h = head;
        // 队列存在, 且 head.waitStatus != 0(即 SIGNAL=-1) 才需要唤醒后继
        if (h != null && h.waitStatus != 0) {
            unparkSuccessor(h);
        }
        return true;
    }
    return false;
}
 
protected final boolean tryRelease(int releases) {
    int c = getState() - releases;       // state--
    if (Thread.currentThread() != getExclusiveOwnerThread())
        throw new IllegalMonitorStateException();
    boolean free = false;
    // 支持重入: state 减到 0 才算真正释放
    if (c == 0) {
        free = true;
        setExclusiveOwnerThread(null);
    }
    setState(c);
    return free;
}
 
private void unparkSuccessor(Node node) {
    int ws = node.waitStatus;
    if (ws < 0) {
        // 把头节点状态从 -1 重置为 0, CAS 失败也没关系
        compareAndSetWaitStatus(node, ws, 0);
    }
    Node s = node.next;
    // 后继不存在或已取消: 从 tail 向前找最靠近 head 的有效节点
    if (s == null || s.waitStatus > 0) {
        s = null;
        for (Node t = tail; t != null && t != node; t = t.prev)
            if (t.waitStatus <= 0)
                s = t;
    }
    if (s != null)
        LockSupport.unpark(s.thread);
}

注意一个分工细节:unparkSuccessor 只负责唤醒,被唤醒节点自己从队列中脱离(抢到锁后 setHead),而不是由释放者替它出队。

可重入原理

可重入完全体现在 state 上,加锁和释放是对称的:

  • 加锁时 nonfairTryAcquire 发现 current == getExclusiveOwnerThread(),说明锁的主人就是自己,不做 CAS,直接 state++;
  • 解锁时 tryRelease 只是 state--,只有 state 减为 0,才清空独占线程、返回 true(真正释放),否则只是减少一层重入计数。
// 加锁: 重入分支
else if (current == getExclusiveOwnerThread()) {
    int nextc = c + acquires;
    if (nextc < 0)
        throw new Error("Maximum lock count exceeded");
    setState(nextc);
    return true;
}
 
// 解锁: 计数归零才释放
int c = getState() - releases;
boolean free = false;
if (c == 0) {
    free = true;
    setExclusiveOwnerThread(null);
}
setState(c);
return free;

所以重入了几层,就必须 unlock 几层。

可打断原理

不可打断模式(lock()):即使等待线程在 park 期间被 interrupt 唤醒,它也不会退出,而是继续驻留在 AQS 队列中循环抢锁,一直等到真正获得锁后,才通过 selfInterrupt() 把中断标记补上:

if (shouldParkAfterFailedAcquire(p, node) &&
    parkAndCheckInterrupt()) {
    // 因 interrupt 被唤醒: 只把 interrupted 标记为 true
    interrupted = true;
}
 
// 拿到锁返回后
if (acquireQueued(addWaiter(Node.EXCLUSIVE), arg)) {
    selfInterrupt();   // 重新产生一次中断, 交给上层处理
}
 
static void selfInterrupt() {
    Thread.currentThread().interrupt();
}

可打断模式(lockInterruptibly())走 acquireInterruptibly:park 过程中一旦被 interrupt,直接抛出 InterruptedException 跳出循环,不再回头抢锁:

public final void acquireInterruptibly(int arg) throws InterruptedException {
    if (Thread.interrupted())
        throw new InterruptedException();
    if (!tryAcquire(arg))
        doAcquireInterruptibly(arg);
}
 
private void doAcquireInterruptibly(int arg) throws InterruptedException {
    final Node node = addWaiter(Node.EXCLUSIVE);
    boolean failed = true;
    try {
        for (;;) {
            final Node p = node.predecessor();
            if (p == head && tryAcquire(arg)) {
                setHead(node);
                p.next = null;
                failed = false;
                return;
            }
            if (shouldParkAfterFailedAcquire(p, node) &&
                parkAndCheckInterrupt()) {
                // park 期间被打断: 直接抛异常, 不会再次进入 for(;;)
                throw new InterruptedException();
            }
        }
    } finally {
        if (failed)
            cancelAcquire(node);
    }
}

公平锁实现原理

公平锁与非公平锁的区别只在 tryAcquire:state 为 0 时,先调用 hasQueuedPredecessors() 检查 AQS 队列里有没有排在自己前面的线程,没有才去 CAS:

static final class FairSync extends Sync {
 
    final void lock() {
        acquire(1);   // 注意: 没有非公平锁开始那次直接 CAS 插队
    }
 
    protected final boolean tryAcquire(int acquires) {
        final Thread current = Thread.currentThread();
        int c = getState();
        if (c == 0) {
            // 先检查队列中是否有前驱节点, 没有才去竞争
            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;
        Node h = head;
        Node s;
        // h != t 说明队列中有节点
        return h != t &&
            ((s = h.next) == null ||        // 老二还没接上(正在入队)
             s.thread != Thread.currentThread()); // 或老二不是自己
    }
}

公平的代价是每次释放都必须按队列顺序唤醒、不能容忍新来的线程插队,线程切换更频繁,所以吞吐量通常低于非公平锁。

条件变量实现原理

每个条件变量对应一个独立的等待队列,实现类是 AQS 的内部类 ConditionObject,它持有 firstWaiter、lastWaiter 两个指针,队列是单向的,节点状态为 -2(Node.CONDITION)。

public class ConditionObject implements Condition, java.io.Serializable {
    private transient Node firstWaiter;   // 第一个等待节点
    private transient Node lastWaiter;    // 最后一个等待节点
}

await 流程(开始时 Thread-0 持有锁):

  1. addConditionWaiter:创建状态为 -2、关联 Thread-0 的新 Node,加入条件等待队列尾部;
  2. fullyRelease:因为线程可能重入过多次,需要把当前 state 一次性全部释放,并记录释放前的 savedState;
  3. 释放后 unpark AQS 队列中的下一个节点竞争锁,假设没有其他竞争者,Thread-1 获得锁;
  4. 只要自己的节点还没被转移回 AQS 队列,就 park 阻塞。

await:Thread-0 进入 ConditionObject 等待队列,状态为 -2

signal 流程(Thread-1 来唤醒 Thread-0):

  1. signal 要求必须持有锁,否则抛 IllegalMonitorStateException;
  2. doSignal 取出条件等待队列中的第一个节点(Thread-0);
  3. transferForSignal:先 CAS 把它的状态从 -2 改为 0,然后调用 enq 把它转移到 AQS 队列尾部,再把它的前驱节点状态置为 -1(SIGNAL),必要时直接 unpark;
  4. Thread-1 之后执行 unlock,Thread-0 在 AQS 队列中被唤醒,重新参与正常的抢锁流程。

signal:Thread-0 的节点从条件队列转移到 AQS 队列尾部

注意:signal 只是把节点"搬家"到 AQS 队列,并不代表它立刻获得锁,仍要按队列规则竞争。这与 wait/notify 中"被唤醒后仍需重新进入 EntryList 竞争"是同一个道理。

核心源码:

// 入条件等待队列
private Node addConditionWaiter() {
    Node t = lastWaiter;
    // 顺手清理已取消的节点
    if (t != null && t.waitStatus != Node.CONDITION) {
        unlinkCancelledWaiters();
        t = lastWaiter;
    }
    Node node = new Node(Thread.currentThread(), Node.CONDITION);
    if (t == null)
        firstWaiter = node;
    else
        t.nextWaiter = node;
    lastWaiter = node;
    return node;
}
 
// signal: 唤醒一个, 把条件队列第一个节点转移走
private void doSignal(Node first) {
    do {
        if ((firstWaiter = first.nextWaiter) == null)
            lastWaiter = null;
        first.nextWaiter = null;
    } while (!transferForSignal(first) &&
             (first = firstWaiter) != null);
}
 
// 节点从条件队列转移到 AQS 队列
final boolean transferForSignal(Node node) {
    // 状态已不是 CONDITION, 说明节点被取消了, 转移失败
    if (!compareAndSetWaitStatus(node, Node.CONDITION, 0))
        return false;
 
    Node p = enq(node);   // 加入 AQS 队列尾部, p 是它的前驱
    int ws = p.waitStatus;
    if (ws > 0 ||                              // 前驱被取消
        !compareAndSetWaitStatus(p, ws, Node.SIGNAL)) {
        LockSupport.unpark(node.thread);       // 兜底: 直接唤醒让它重新同步
    }
    return true;
}
 
// 重入可能导致 state > 1, 因此要一次性全部释放
final int fullyRelease(Node node) {
    boolean failed = true;
    try {
        int savedState = getState();
        if (release(savedState)) {
            failed = false;
            return savedState;   // 供以后重新获取时恢复计数
        } else {
            throw new IllegalMonitorStateException();
        }
    } finally {
        if (failed)
            node.waitStatus = Node.CANCELLED;
    }
}

可打断的 await() 主循环如下,唤醒后还要在 AQS 队列中把锁抢回来,才算真正从 await 返回:

public final void await() throws InterruptedException {
    if (Thread.interrupted()) {
        throw new InterruptedException();
    }
    Node node = addConditionWaiter();
    int savedState = fullyRelease(node);
    int interruptMode = 0;
    // 节点还没有转移到 AQS 队列, 就一直阻塞
    while (!isOnSyncQueue(node)) {
        LockSupport.park(this);
        if ((interruptMode = checkInterruptWhileWaiting(node)) != 0)
            break;
    }
    // 转移后还要按 AQS 规则重新获取锁
    if (acquireQueued(node, savedState) && interruptMode != THROW_IE)
        interruptMode = REINTERRUPT;
    if (node.nextWaiter != null)
        unlinkCancelledWaiters();
    if (interruptMode != 0)
        reportInterruptAfterWait(interruptMode);
}

signalAll 的区别只是 doSignalAll:把等待队列中的所有节点依次转移到 AQS 队列。

读写锁 ReentrantReadWriteLock 原理

state 的拆分

读写锁的读锁和写锁用的是同一个 Sync 同步器,因此等待队列、state 等也是同一个。一个 int state 被拆成两半:

  • 低 16 位:写锁计数(独占)
  • 高 16 位:读锁计数(共享,每有一个线程加读锁就 +1,加的是 SHARED_UNIT = 1 << 16)

t1 执行 w.lock(),流程与 ReentrantLock 没有本质区别,只是占用低 16 位,图中记作 0_1:

t1 获取写锁,state 低 16 位为 1

共享获取的返回值 tryAcquireShared 有三种语义:

  • -1:获取失败;
  • 0:获取成功,但后继节点不会继续被唤醒;
  • 正数:获取成功,且还有几个后继需要唤醒(读写锁成功时返回 1)。

图解加锁与释放时序

t2 执行 r.lock:进入 acquireShared(1) → tryAcquireShared 发现有写锁占据,返回 -1;进入 doAcquireShared,addWaiter 时节点模式是 Node.SHARED(不是 EXCLUSIVE)。t2 先看自己是不是老二,是则再试一次;失败后把前驱 waitStatus 改为 -1,再试一次仍失败,于是 park。

t3 r.lock、t4 w.lock:同样入队,形成共享、共享、独占的混合队列,三个节点的前驱状态都被置为 -1:

t2、t3 读锁与 t4 写锁排队:Shared、Shared、Exclusive 混合队列

t1 w.unlock:走写锁的 release(1),tryRelease 成功后执行 unparkSuccessor,老二 t2 在 parkAndCheckInterrupt() 处恢复,再来一轮 for 循环执行 tryAcquireShared 成功,读锁计数加一。

接着 t2 调用 setHeadAndPropagate(node, 1) 把自己设为头节点——共享模式的关键在这里:它还会检查下一个节点是不是 shared,如果是,调用 doReleaseShared() 把 head 状态从 -1 改回 0 并继续唤醒老二,于是 t3 也恢复运行,tryAcquireShared 成功,读锁计数再加一。两个读锁可以同时持有,这就是"读-读不互斥",此时 state 高 16 位为 2:

t2、t3 共享传播:两个读锁同时获得,state 高 16 位为 2

t3 成为 head 后再看后继:t4 是独占节点,不是 shared,传播到此为止,t4 继续 park——读锁持有时写锁必须等待,避免写饥饿靠的是后面会讲到的 readerShouldBlock。

t2 r.unlock、t3 r.unlock:都走 releaseShared(1) → tryReleaseShared,CAS 让高 16 位计数减一。t2 释放后计数仍为 1,不唤醒任何人;t3 释放后计数归零,tryReleaseShared 返回 true,进入 doReleaseShared() 将头节点状态从 -1 改 0 并唤醒老二。t4 在 acquireQueued(独占节点走的仍是独占流程)中恢复,自己已是老二且没有竞争,tryAcquire 成功,修改头节点,流程结束:

最后一个读锁释放,计数归零,t4 被唤醒获得写锁

写锁加锁与释放源码

// WriteLock.lock
public void lock() {
    sync.acquire(1);
}
 
protected final boolean tryAcquire(int acquires) {
    Thread current = Thread.currentThread();
    int c = getState();
    int w = exclusiveCount(c);   // 低 16 位, 写锁计数
 
    if (c != 0) {
        if (
            // c != 0 但 w == 0: 有读锁存在, 写锁不能获取
            w == 0 ||
            // 或者写锁的主人不是自己
            current != getExclusiveOwnerThread()
        ) {
            return false;
        }
        if (w + exclusiveCount(acquires) > MAX_COUNT)
            throw new Error("Maximum lock count exceeded");
        // 写锁重入
        setState(c + acquires);
        return true;
    }
    if (
        // 公平锁下该排队就排队; 非公平锁 writerShouldBlock 恒为 false
        writerShouldBlock() ||
        !compareAndSetState(c, c + acquires)
    ) {
        return false;
    }
    setExclusiveOwnerThread(current);
    return true;
}
 
// WriteLock.unlock
protected final boolean tryRelease(int releases) {
    if (!isHeldExclusively())
        throw new IllegalMonitorStateException();
    int nextc = getState() - releases;
    // 写锁计数(低16位)减为 0 才算释放成功
    boolean free = exclusiveCount(nextc) == 0;
    if (free) {
        setExclusiveOwnerThread(null);
    }
    setState(nextc);
    return free;
}

读锁加锁与释放源码

// ReadLock.lock
public void lock() {
    sync.acquireShared(1);
}
 
public final void acquireShared(int arg) {
    if (tryAcquireShared(arg) < 0) {
        doAcquireShared(arg);
    }
}
 
protected final int tryAcquireShared(int unused) {
    Thread current = Thread.currentThread();
    int c = getState();
    // 其它线程持有写锁, 读锁失败
    // 注意: 写锁主人是自己时允许再加读锁(锁降级)
    if (exclusiveCount(c) != 0 &&
        getExclusiveOwnerThread() != current) {
        return -1;
    }
    int r = sharedCount(c);   // 高 16 位读锁计数
    if (
        // 非公平实现: 队列第一个是写锁时, 读锁该阻塞(防止写线程被源源不断的读饿死)
        !readerShouldBlock() &&
        r < MAX_COUNT &&
        // 高 16 位 +1, 即 state + SHARED_UNIT
        compareAndSetState(c, c + SHARED_UNIT)
    ) {
        // ... 维护 firstReader / HoldCounter 等, 省略
        return 1;
    }
    // 快速尝试失败(可能有 CAS 竞争), 进入自旋完整版
    return fullTryAcquireShared(current);
}
 
// 非公平锁的 readerShouldBlock: 队列头一个节点是独占(写)节点就阻塞读
final boolean readerShouldBlock() {
    return apparentlyFirstQueuedIsExclusive();
}

doAcquireShared 与独占版的区别只在成功后的动作——setHeadAndPropagate:

private void doAcquireShared(int arg) {
    final Node node = addWaiter(Node.SHARED);   // 共享节点
    boolean failed = true;
    try {
        boolean interrupted = false;
        for (;;) {
            final Node p = node.predecessor();
            if (p == head) {
                int r = tryAcquireShared(arg);
                if (r >= 0) {
                    // r > 0 表示允许向后传播唤醒
                    setHeadAndPropagate(node, r);
                    p.next = null;
                    if (interrupted)
                        selfInterrupt();
                    failed = false;
                    return;
                }
            }
            if (shouldParkAfterFailedAcquire(p, node) &&
                parkAndCheckInterrupt()) {
                interrupted = true;
            }
        }
    } finally {
        if (failed)
            cancelAcquire(node);
    }
}
 
private void setHeadAndPropagate(Node node, int propagate) {
    Node h = head;          // 记录旧 head
    setHead(node);
    // propagate > 0, 或新旧 head 状态为 SIGNAL/PROPAGATE, 都尝试继续释放
    if (propagate > 0 || h == null || h.waitStatus < 0 ||
        (h = head) == null || h.waitStatus < 0) {
        Node s = node.next;
        // 后继不存在, 或后继是共享节点: 继续唤醒
        if (s == null || s.isShared()) {
            doReleaseShared();
        }
    }
}
 
// 读锁释放: 计数减一, 只有归零才真正触发唤醒
protected final boolean tryReleaseShared(int unused) {
    for (;;) {
        int c = getState();
        int nextc = c - SHARED_UNIT;
        if (compareAndSetState(c, nextc)) {
            // 读计数不影响其它读线程, 但影响等写锁的线程
            return nextc == 0;
        }
    }
}
 
public final boolean releaseShared(int arg) {
    if (tryReleaseShared(arg)) {
        doReleaseShared();
        return true;
    }
    return false;
}

doReleaseShared 是共享模式释放/传播的核心,它在多个释放线程并发执行时靠 CAS 保证只唤醒一次,并且把 head 的 0 状态进一步置成 PROPAGATE(-3),这就是下一节要讲的 bug 修复:

private void doReleaseShared() {
    for (;;) {
        Node h = head;
        if (h != null && h != tail) {
            int ws = h.waitStatus;
            if (ws == Node.SIGNAL) {
                // SIGNAL -> 0 成功才唤醒, CAS 失败说明别的线程在释放, 重试复查
                if (!compareAndSetWaitStatus(h, Node.SIGNAL, 0))
                    continue;
                unparkSuccessor(h);
            }
            // 已经是 0: 置为 PROPAGATE(-3), 保证唤醒信号能继续传播
            else if (ws == 0 &&
                     !compareAndSetWaitStatus(h, 0, Node.PROPAGATE))
                continue;
        }
        if (h == head)
            break;
    }
}

Semaphore 原理

加锁解锁流程

Semaphore 像一个停车场:permits 就是停车位数量,线程获取许可相当于抢到车位,空余车位减一;归还许可后车位加一。

刚开始 permits(state)为 3,5 个线程同时来获取。假设 Thread-1、Thread-2、Thread-4 CAS 竞争成功,Thread-0 和 Thread-3 竞争失败,进入 AQS 队列 park(共享节点)。之后 Thread-4 释放一个许可,state 回到 1,Thread-0 被唤醒后竞争成功,许可再次归零,它成为新 head 并 unpark 后继 Thread-3;但此时 permits 是 0,Thread-3 尝试获取失败后会再次 park:

Semaphore:三个线程持有许可,Thread-0、Thread-3 在共享队列中等待

源码分析

Semaphore 的 Sync 同样分公平/非公平,acquire 走共享的可打断获取,state 存的就是剩余许可数:

static final class NonfairSync extends Sync {
 
    NonfairSync(int permits) {
        super(permits);   // permits 即 state
    }
 
    public void acquire() throws InterruptedException {
        sync.acquireSharedInterruptibly(1);
    }
 
    public final void acquireSharedInterruptibly(int arg)
            throws InterruptedException {
        if (Thread.interrupted())
            throw new InterruptedException();
        if (tryAcquireShared(arg) < 0)
            doAcquireSharedInterruptibly(arg);
    }
 
    final int nonfairTryAcquireShared(int acquires) {
        for (;;) {
            int available = getState();
            int remaining = available - acquires;
            if (
                // 许可已用完: 返回负数, 进入队列阻塞
                remaining < 0 ||
                // CAS 成功: 返回剩余许可数(可能是 0)
                compareAndSetState(available, remaining)
            ) {
                return remaining;
            }
        }
    }
 
    public void release() {
        sync.releaseShared(1);
    }
 
    protected final boolean tryReleaseShared(int releases) {
        for (;;) {
            int current = getState();
            int next = current + releases;
            if (next < current)
                throw new Error("Maximum permit count exceeded");
            if (compareAndSetState(current, next))
                return true;
        }
    }
}

doAcquireSharedInterruptibly 与读写锁的 doAcquireShared 几乎一样:被唤醒后调用 setHeadAndPropagate(node, r),注意 r 是剩余资源数——为 0 时不会继续向后继传播。

为什么要有 PROPAGATE

PROPAGATE(waitStatus = -3)是为修复共享模式下的一个唤醒信号丢失 bug 而引入的。先看修复前的实现:

// 修复前的 releaseShared
public final boolean releaseShared(int arg) {
    if (tryReleaseShared(arg)) {
        Node h = head;
        if (h != null && h.waitStatus != 0)
            unparkSuccessor(h);
        return true;
    }
    return false;
}
 
// 修复前的 setHeadAndPropagate
private void setHeadAndPropagate(Node node, int propagate) {
    setHead(node);
    if (propagate > 0 && node.waitStatus != 0) {
        Node s = node.next;
        if (s == null || s.isShared())
            unparkSuccessor(node);
    }
}

假设队列状态为 head(-1) -> t1(-1) -> t2(0),t1、t2 都在 park,t3、t4 两个线程将要释放许可,顺序是先 t3 后 t4。出 bug 的时序是:

  1. t3 调用 releaseShared(1),直接 unparkSuccessor(head),head 的 waitStatus 从 -1 变为 0;
  2. t1 被唤醒,执行 tryAcquireShared,返回 0——获取成功,但没有剩余资源;
  3. t4 接着调用 releaseShared(1),此时它读到的 head.waitStatus 是 0(和第 1 步是同一个 head),不满足 waitStatus != 0 的条件,不调用 unparkSuccessor;
  4. t1 获取成功执行 setHeadAndPropagate,因为 propagate == 0 不满足 propagate > 0,不唤醒后继。结果 t2 明明应该被第二个许可唤醒,却永远睡死在队列里。

问题的根源是:t4 释放的信号恰好落在"head 状态为 0、t1 还没完成出队"的空档里,信号被丢了。

修复思路是增加一个中间状态 PROPAGATE(-3):t4 释放时即使 head 已经是 0,也通过 doReleaseShared 把它 CAS 成 -3,表示"有一次释放正在传播";t1 出队时 setHeadAndPropagate 发现旧 head 的 waitStatus < 0(-1 或 -3 都算),就会继续执行 doReleaseShared() 唤醒 t2。修复后时序:

  1. t3 releaseShared() → unparkSuccessor(head),head 状态 -1 → 0;
  2. t1 被唤醒,tryAcquireShared 返回 0;
  3. t4 releaseShared(),此时 head 为 0,调用 doReleaseShared() 把状态置为 PROPAGATE(-3);
  4. t1 获取成功调用 setHeadAndPropagate 时读到 h.waitStatus < 0,于是调用 doReleaseShared(),t2 被正确唤醒。

修复后的两个方法就是前面读写锁一节看到的版本:判断条件从单纯的 propagate > 0 扩展为同时检查新旧 head 的 waitStatus < 0,释放动作统一走会循环重试、会把 0 置为 -3 的 doReleaseShared。读写锁能连续唤醒多个排队的读线程,靠的也是这套机制。

CountDownLatch 原理

CountDownLatch 同样是 AQS 共享模式的"薄封装":构造时给定计数,写入 AQS 的 state;await() 让线程在计数未归零时排队阻塞,countDown() 让计数减一,计数归零时一次性唤醒所有等待线程。

它的 Sync 只需要实现两个钩子:

private static final class Sync extends AbstractQueuedSynchronizer {
 
    Sync(int count) {
        setState(count);
    }
 
    int getCount() {
        return getState();
    }
 
    // await 走 acquireShared: state 不为 0 就返回 -1(排队), 为 0 返回 1(放行)
    protected int tryAcquireShared(int acquires) {
        return (getState() == 0) ? 1 : -1;
    }
 
    // countDown 走 releaseShared: CAS 把 state 减 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;     // 只有减到 0 的那个线程负责唤醒
        }
    }
}
 
public void await() throws InterruptedException {
    sync.acquireSharedInterruptibly(1);
}
 
public void countDown() {
    sync.releaseShared(1);
}

唤醒之所以是"一次性广播",正是因为共享模式的传播机制:计数归零的那次 tryReleaseShared 返回 true,进入 doReleaseShared 唤醒队首节点;被唤醒的节点 tryAcquireShared 返回正数 1,setHeadAndPropagate 又会继续唤醒后面的 shared 节点,如此传播,所有在 await() 上等待的线程被依次放行。这也是 CountDownLatch 一旦归零就不能复用、而 CyclicBarrier 可以重置的根本原因——它没有任何机制把 state 再加回去。

小结

  • AQS 用 volatile + CAS 的 state 表达资源状态,用 CLH 变体的 FIFO 双向队列安置失败线程,用 park/unpark 做阻塞与唤醒,子类只需实现 tryAcquire/tryRelease 等几个钩子。
  • ReentrantLock 的非公平锁靠"入队前先 CAS 抢一次、唤醒时允许新来者插队"换取吞吐;公平锁靠 hasQueuedPredecessors 保证顺序;可重入是 state 的加减;可打断分"获锁后补中断"和"直接抛异常"两种;ConditionObject 维护独立的单向等待队列,await 时释放锁并 park,signal 时把节点搬回 AQS 队列重新竞争。
  • 读写锁把 state 拆成高低 16 位,共享节点成功后通过 setHeadAndPropagate + doReleaseShared 把唤醒传播给后续读节点,遇到独占节点停止;readerShouldBlock 防止写线程被读流量饿死。
  • Semaphore 用共享模式管理许可;PROPAGATE(-3) 状态修复了并发释放时唤醒信号在 head 状态切换空档中丢失的 bug。
  • CountDownLatch 用共享获取的"非 0 即失败"和归零后的传播唤醒,实现了一个计数到 0 的一次性门闩。