Skip to content

AQS

AbstractQueuedSynchronizer简称AQS是一个抽象类,AQS定义了一套多线程访问共享资源的同步器框架,是一个依赖状态(state)的同步器。

JUC(java.util.concurrent包)下面的Lock和其他一些并发工具类都是基于它来实现的。AQS维护了一个volatile的state和一个CLH(FIFO)双向队列。

如果你是第一次学习 AQS,建议先看课程型主线:AQS 与 JUC 工具类全过程原理。那篇会先讲清 state、同步队列、park/unpark、独占/共享模式、ReentrantLock、Semaphore、CountDownLatch 的完整流程,再回到本页看源码片段会更顺。

零基础可以先这样理解:

AQS = 一个“排队拿资源”的通用骨架。

很多并发工具的共同问题都是:资源只有一份或有限几份,线程来了以后,能拿到资源就继续执行,拿不到就排队休眠,等前面的线程释放资源后再被唤醒。AQS把“排队、阻塞、唤醒、状态维护”这些通用逻辑做好,具体工具只需要定义“怎样算拿到资源”和“怎样算释放资源”。

为什么 JUC 要抽出 AQS

如果每个同步工具都自己实现排队、阻塞、唤醒、中断、超时和状态修改,代码会大量重复,而且极容易写错。AQS 把这些通用能力沉淀成模板,ReentrantLockSemaphoreCountDownLatch 等工具只需要围绕 state 定义自己的获取和释放规则。

工具state 可以理解成获取规则
ReentrantLock是否被占用以及重入次数没人持有或当前线程重入
Semaphore剩余许可证数量许可证数量大于 0
CountDownLatch还剩多少次倒计数state 归零后等待线程通过

如果不会 AQS,学习 JUC 时就只能背 API;理解 AQS 后,可以看懂这些工具为什么行为不同,但底层都围绕队列和状态流转。

工作原理图

mermaid
flowchart TD
    A["线程调用 lock/acquire"] --> B{"tryAcquire 是否成功"}
    B -- "成功" --> C["修改 state\n设置持有线程\n继续执行业务代码"]
    B -- "失败" --> D["封装为 Node"]
    D --> E["CAS 加入 CLH 等待队列尾部"]
    E --> F{"前驱是否为 head\n且再次抢锁成功"}
    F -- "是" --> G["当前节点成为 head\n线程继续执行"]
    F -- "否" --> H["LockSupport.park 挂起线程"]
    I["持锁线程 unlock/release"] --> J["tryRelease 修改 state"]
    J --> K{"资源是否完全释放"}
    K -- "是" --> L["unparkSuccessor 唤醒后继节点"]
    L --> F
    K -- "否" --> M["可重入层数未归零\n不唤醒后继"]

这张图抓住 AQS 的核心:state 表示资源状态,CLH 队列保存等待线程,CAS 保证并发修改安全,park/unpark 负责让线程休眠和唤醒。

同步状态-state变量

state是由volatile修饰的int类型,用来表示当前线程同步状态。

java
// volatile修饰的state
private volatile int state;
protected final int getState() {
    return state;
}
protected final void setState(int newState) {
    state = newState;
}
/**
 * 使用CAS+volatile,基于原子性与可见性的对state进行设值
 * expect: 期望值
 * update: 更新值
 */
protected final boolean compareAndSetState(int expect, int update) {
    // 使用Unsafe类,调用JNI方法
    return unsafe.compareAndSwapInt(this, stateOffset, expect, update);
}

CLH(FIFO)队列

AQS中是通过内部类Node来维护一个CLH队列的。

CLH队列,全称Craig-Landin-Hagersten队列,是一种基于链表结构的自旋锁等待队列。

源码如下:

java
    static final class Node {
        /** 标记共享式访问 */
        static final Node SHARED = new Node();
        /** 标记独占式访问 */
        static final Node EXCLUSIVE = null;

        /** 字段waitStatus的值,表示当前节点已取消等待 */
        static final int CANCELLED =  1;
        /** 字段waitStatus的值,表示当前节点取消或释放资源后,通知下一个节点 */
        static final int SIGNAL    = -1;
        /** 表示正在等待触发条件 */
        static final int CONDITION = -2;
        /**
         * 表示下一个共享获取应无条件传播
         */
        static final int PROPAGATE = -3;

        /**
         * Status field, taking on only the values:
         *   SIGNAL:     此节点的后续节点被(或很快将被)阻止(通过park),
         *               因此当前节点在释放或取消时必须取消标记其后续节点。
         *               为了避免竞争,获取方法必须首先指示它们需要一个信号,
         *               然后重试原子获取,然后在失败时阻止。
         *   CANCELLED:  由于超时或中断,此节点被取消。节点永远不会离开此状态。
         *               特别是,具有取消节点的线程再也不会阻塞。
         *   CONDITION:  此节点当前在条件队列中。
         *               在传输之前,它不会用作同步队列节点,此时状态将设置为0。
         *               (此处使用此值与字段的其他用途无关,但简化了机制。)
         *   PROPAGATE:  releaseShared应该传播到其他节点。
         *               这是在doReleaseShared中设置的(仅适用于头节点),
         *               以确保传播继续进行,即使其他操作已经介入。
         *   0:          以上都没有
         *
         * 这些值以数字形式排列,以简化使用。非负值表示节点不需要发出信号。所以,大多数代码不需要检查特定的值,只需要检查符号。
         *
         * 对于正常同步节点,该字段初始化为0,对于条件节点,则初始化为CONDITION。它使用CAS(或在可能的情况下,无条件的易失性写入)进行修改。
         */
        volatile int waitStatus;

        /**
         * 前一个节点
         */
        volatile Node prev;

        /**
         * 下一个节点
         */
        volatile Node next;

        /**
         * 将此节点排入队列的线程。在构造时初始化,使用后为null。
         * 节点对应线程
         */
        volatile Thread thread;

        /**
         * 下一个等待的节点
         */
        Node nextWaiter;

        /**
         * 是否是共享式访问
         */
        final boolean isShared() {
            return nextWaiter == SHARED;
        }

        /**
         * 返回上一个节点,如果为null,则抛出NullPointerException。
         * 前置任务不能为null时使用。可以取消空检查,但存在空检查是为了帮助VM。
         */
        final Node predecessor() {
            Node p = prev;
            if (p == null)
                throw new NullPointerException();
            else
                return p;
        }

        /** 共享式访问的构造函数 */
        Node() {}

        /** 添加下一个等待者的构造 */
        Node(Node nextWaiter) {
            this.nextWaiter = nextWaiter;
            THREAD.set(this, Thread.currentThread());
        }

        /** addConditionWaiter使用的构造函数。 */
        Node(int waitStatus) {
            WAITSTATUS.set(this, waitStatus);
            THREAD.set(this, Thread.currentThread());
        }

        /** CASes waitStatus field. */
        final boolean compareAndSetWaitStatus(int expect, int update) {
            return WAITSTATUS.compareAndSet(this, expect, update);
        }

        /** CASes next field. */
        final boolean compareAndSetNext(Node expect, Node update) {
            return NEXT.compareAndSet(this, expect, update);
        }

        final void setPrevRelaxed(Node p) {
            PREV.set(this, p);
        }

        // VarHandle mechanics
        private static final VarHandle NEXT;
        private static final VarHandle PREV;
        private static final VarHandle THREAD;
        private static final VarHandle WAITSTATUS;
        static {
            try {
                MethodHandles.Lookup l = MethodHandles.lookup();
                NEXT = l.findVarHandle(Node.class, "next", Node.class);
                PREV = l.findVarHandle(Node.class, "prev", Node.class);
                THREAD = l.findVarHandle(Node.class, "thread", Thread.class);
                WAITSTATUS = l.findVarHandle(Node.class, "waitStatus", int.class);
            } catch (ReflectiveOperationException e) {
                throw new ExceptionInInitializerError(e);
            }
        }
    }

临界资源获取

两种获取方式

  • 独占式(EXCLUSIVE)

仅有一个线程能在同一时刻获取到资源并处理,如ReentrantLock的实现。

  • 共享式(SHARED)

多个线程可以同时获取到资源并处理,如Semaphore/CountDownLatch等。

AQS中大部分逻辑已经被实现,集成类只需要重写state的获取(acquire)与释放(release)方法,因为在AQS中,这些方法默认定义的实现方式都是抛出不支持操作异常,所以按需实现即可。 其中需要继承类重写的方法有:

  • tryAcquire(int arg)

此方法是独占式的获取资源方法,成功则返回true,失败返回false。

  • tryRelease(int arg)

此方法是独占式的释放资源方法,成功则返回true,失败返回false。

  • tryAcquireShared(int arg)

此方法是共享式的获取资源方法,返回负数表示失败,0表示获取成功,但是没有可用资源,正数表示获取成功,且有可用资源。

  • tryReleaseShared(int arg)

此方法是共享式的释放资源方法,如果允许唤醒后续等待线程则返回true,不允许则返回false。

  • isHeldExclusively()

判断当前线程是否正在独享资源,是则返回true,否则返回false。

最小 Demo:用 AQS 实现一个不可重入互斥锁

下面这个例子不是为了替代 ReentrantLock,而是帮助理解 AQS 的扩展点。它只允许一个线程持有锁,其他线程必须排队。

java
import java.util.concurrent.locks.AbstractQueuedSynchronizer;
import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.Condition;
import java.util.concurrent.TimeUnit;

public class SimpleMutex implements Lock {
    private static class Sync extends AbstractQueuedSynchronizer {
        @Override
        protected boolean tryAcquire(int arg) {
            if (compareAndSetState(0, 1)) {
                setExclusiveOwnerThread(Thread.currentThread());
                return true;
            }
            return false;
        }

        @Override
        protected boolean tryRelease(int arg) {
            if (getState() == 0) {
                throw new IllegalMonitorStateException();
            }
            setExclusiveOwnerThread(null);
            setState(0);
            return true;
        }

        @Override
        protected boolean isHeldExclusively() {
            return getState() == 1 && getExclusiveOwnerThread() == Thread.currentThread();
        }

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

    private final Sync sync = new Sync();

    @Override
    public void lock() {
        sync.acquire(1);
    }

    @Override
    public void unlock() {
        sync.release(1);
    }

    @Override
    public boolean tryLock() {
        return sync.tryAcquire(1);
    }

    @Override
    public void lockInterruptibly() throws InterruptedException {
        sync.acquireInterruptibly(1);
    }

    @Override
    public boolean tryLock(long time, TimeUnit unit) throws InterruptedException {
        return sync.tryAcquireNanos(1, unit.toNanos(time));
    }

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

这个 Demo 对应到 AQS 的职责划分:

位置做什么
tryAcquire定义什么时候能拿到锁,这里是 state 从 0 改成 1
tryRelease定义什么时候释放锁,这里是把 state 改回 0
acquire/releaseAQS 已经实现,负责排队、阻塞、唤醒
ConditionObjectAQS 提供的条件队列实现

独占模式-获取资源

使用AQS中的acquire(int arg)方法

java
/**
 * 以独占模式获取,忽略中断。通过至少调用一次tryAcquire来实现,成功后返回。
 * 否则,线程将排队,可能会重复阻塞和取消阻塞,调用tryAcquire直到成功。
 * 此方法可用于实现方法Lock.lock。
 */
public final void acquire(int arg) {
    if (!tryAcquire(arg) &&
        acquireQueued(addWaiter(Node.EXCLUSIVE), arg))
        selfInterrupt();
}

该方法分为4个部分:

  • tryAcquire

需要自己实现的方法,如果获取到资源使用权,则返回true,反之fasle。如果获取到资源,返回true,!true为false,根据&&的短路性,则不会执行后续方法,直接跳过程序。如果未获取到资源,返回false,!false为true,则进入后续方法。

  • addWaiter 如果未获取到资源使用权,则首先会调用此方法。源码:
java
/**
 * 为当前线程和给定模式创建节点并将其排入队列。
 */
private Node addWaiter(Node mode) {
    // 封装当前线程和独占模式
    Node node = new Node(Thread.currentThread(), mode);
    // 获取尾部节点    
    Node pred = tail;
    if (pred != null) {
        node.prev = pred;
        // CAS设置尾部节点
        if (compareAndSetTail(pred, node)) {
            pred.next = node;
            return node;
        }
    }
    // 如果尾结点为空或者设置尾结点失败
    enq(node);
    return node;
}

/**
 * 将节点插入队列,必要时进行初始化。
 */
private Node enq(final Node node) {
    // 如果CAS设置未成功则死循环
    for (;;) {
        // 获得尾结点
        Node t = tail;
        // 如果尾节点为空,说明CLH队列为空,需要初始化
        if (t == null) {
            if (compareAndSetHead(new Node()))
                 tail = head;
        } else {
            // 设置当前节点的前驱节点
            node.prev = t;
            // CAS设置当前节点为尾结点
            if (compareAndSetTail(t, node)) {
                t.next = node;
                return t;
            }
        }
    }
}
  • acquiredQueued
java
/**
 * 以独占不间断模式为队列中已存在的线程获取。
 * 由条件等待方法和获取方法使用。
 */
final boolean acquireQueued(final Node node, int arg) {
    // 标识资源获取是否失败
    boolean failed = true;
    try {
        // 标识线程是否中断
        boolean interrupted = false;
        for (;;) {
            // 获得当前节点的前驱节点
            final Node p = node.predecessor();
            // 如果前驱节点为头结点,说明快到当前节点了,尝试获取资源
            if (p == head && tryAcquire(arg)) {
                // 获取资源成功
                // 设置当前节点为头结点
                setHead(node);
                // 取消前驱节点(以前的头部)的后节点,方便GC回收
                p.next = null; // help GC
                // 标识未失败
                failed = false;
                // 返回中断标志
                return interrupted;
            }
            // 如果当前节点的前驱节点不是头结点或获取资源失败
            // 需要用shouldParkAfterFailedAcquire函数判断是否需要阻塞该节点持有的线程
            // 如果需要阻塞,则执行parkAndCheckInterrupt方法,并设置被中断
            if (shouldParkAfterFailedAcquire(p, node) &&
                parkAndCheckInterrupt())
                interrupted = true;
        }
    } finally {
        // 如果最终获取资源失败
        if (failed)
            // 当前节点取消获取资源
            cancelAcquire(node);
    }
}
  • selfInterrupt

中断当前线程

java
static void selfInterrupt() {
    Thread.currentThread().interrupt();
}

独占模式-释放资源

  • release

释放资源并唤醒后继线程

java
/**
 * 以独占模式发布。如果tryRelease返回true,则通过取消阻止一个或多个线程来实现。此方法可用于实现方法Lock.unlock
 *
 * @param arg the release argument.  This value is conveyed to
 *        {@link #tryRelease} but is otherwise uninterpreted and
 *        can represent anything you like.
 * @return the value returned from {@link #tryRelease}
 */
public final boolean release(int arg) {
    if (tryRelease(arg)) {
        // 获取头结点
        Node h = head;
        // 头结点不为空且等待状态值不为0
        if (h != null && h.waitStatus != 0)
            // 唤醒后续等待线程
            unparkSuccessor(h);
        return true;
    }
    return false;
}
  • tryRelease

和tryAcquire方法同理需要自己实现的方法,

  • unparkSuccessor

唤醒节点的后续节点-等待线程(如果存在)

java
private void unparkSuccessor(Node node) {
    /*
     * 如果状态为负(即,可能需要信号),则尝试清除信号。如果失败或等待线程更改了状态,则可以。
     */
    int ws = node.waitStatus;
    // 如果等待状态值小于0
    if (ws < 0)
        // 使用CAS将waitStatus设置为0
        compareAndSetWaitStatus(node, ws, 0);

    /*
     * 要取消标记的线程保存在后续节点中,后者通常只是下一个节点。但如果取消或明显为空,则从尾部向后遍历,以找到实际的未取消的后续项。
     */
    Node s = node.next;
    // 如果当前节点没有后继节点或者后继节点放弃竞争资源
    if (s == null || s.waitStatus > 0) {
        s = null;
        // 从队列尾部循环直到当前节点,找到最近的且等待状态值小于0的节点
        for (Node t = tail; t != null && t != node; t = t.prev)
            if (t.waitStatus <= 0)
                s = t;
    }
    // 如果找到的后继节点不为空,则唤醒其持有的线程
    if (s != null)
        LockSupport.unpark(s.thread);
}