Administrator
发布于 2019-05-24 / 1535 阅读
41

AQS 源码初探:ReentrantLock 是怎么实现的

jstack 里那个等了 40 秒的线程

5 月中旬,运营反馈大客户导出特别慢,一个 20 万行的订单导出要跑三分钟。我抓了几份 jstack,发现有个线程的行为很奇怪:

"export-thread-7" #48 prio=5 os_prio=0 tid=0x00007f8c4c0d8000 nid=0x5d43 waiting on condition
   java.lang.Thread.State: WAITING (parking)
    at sun.misc.Unsafe.park(Native Method)
    - parking to wait for  <0x00000006c0a1b3e0>
      (a java.util.concurrent.locks.ReentrantLock$NonfairSync)
    at java.util.concurrent.locks.LockSupport.park(LockSupport.java:175)
    at java.util.concurrent.locks.AbstractQueuedSynchronizer.parkAndCheckInterrupt(AbstractQueuedSynchronizer.java:836)
    at java.util.concurrent.locks.AbstractQueuedSynchronizer.acquireQueued(AbstractQueuedSynchronizer.java:870)
    at java.util.concurrent.locks.AbstractQueuedSynchronizer.acquire(AbstractQueuedSynchronizer.java:1199)
    at java.util.concurrent.locks.ReentrantLock$NonfairSync.lock(ReentrantLock.java:209)
    at java.util.concurrent.locks.ReentrantLock.lock(ReentrantLock.java:285)
    at com.xxx.export.ExcelWriter.writeRow(ExcelWriter.java:63)

连着抓了五次,export-thread-7 全都在 parking。加了埋点统计每个线程累计等待锁的时间,结果是:8 个导出线程里,等待最长的一个累计等了 41.3 秒,最短的一个只等了 2.1 秒。

锁是这么用的,非公平锁(new ReentrantLock() 默认就是非公平的):

@Component
public class ExcelWriter {
    // POI 的 XSSFRow 不是线程安全的,加锁保护
    private final ReentrantLock lock = new ReentrantLock();

    public void writeRow(Sheet sheet, OrderVO vo) {
        lock.lock();
        try {
            Row row = sheet.createRow(sheet.getLastRowNum() + 1);
            row.createCell(0).setCellValue(vo.getOrderNo());
            row.createCell(1).setCellValue(vo.getAmount().doubleValue());
            // ... 共 14 列
        } finally {
            lock.unlock();
        }
    }
}

为了搞明白为什么会这样,我把 AbstractQueuedSynchronizer 的源码读了一遍。

AQS 的三件套

AQS 是 ReentrantLockSemaphoreCountDownLatchReentrantReadWriteLock 这些同步器的公共底座。它把"排队、阻塞、唤醒"这些通用的活干完,把"什么条件下算获取成功"留给子类实现。

它只有三样东西:

  • 一个 volatile 的 state。在 ReentrantLock 里表示重入次数,0 表示没人占用;在 Semaphore 里表示剩余许可数;在 CountDownLatch 里表示计数。
  • 一个 FIFO 的等待队列(变种的 CLH 队列)。双向链表,头节点是哑节点,代表当前持有锁的线程。
  • 一组模板方法acquirerelease 这些是 final 的,子类只需要实现 tryAcquiretryRelease
public abstract class AbstractQueuedSynchronizer
        extends AbstractOwnableSynchronizer {

    private volatile int state;
    private transient volatile Node head;
    private transient volatile Node tail;

    protected final int getState() { return state; }
    protected final void setState(int newState) { state = newState; }
    protected final boolean compareAndSetState(int expect, int update) {
        return unsafe.compareAndSwapInt(this, stateOffset, expect, update);
    }

    // 留给子类实现
    protected boolean tryAcquire(int arg) {
        throw new UnsupportedOperationException();
    }
}

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;   // 在 Condition 队列里
    static final int PROPAGATE  = -3;   // 共享模式下传播唤醒

    volatile int waitStatus;
    volatile Node prev;
    volatile Node next;
    volatile Thread thread;
}

加锁流程走一遍

入口是 ReentrantLock.lock(),它委派给内部类 Sync。先看非公平版本:

// NonfairSync
static final class NonfairSync extends Sync {
    final void lock() {
        if (compareAndSetState(0, 1))          // 第一次抢:直接 CAS,不看队列
            setExclusiveOwnerThread(Thread.currentThread());
        else
            acquire(1);                        // 抢失败,进 AQS 的标准流程
    }

    protected final boolean tryAcquire(int acquires) {
        return nonfairTryAcquire(acquires);
    }
}

// 父类 Sync 里的实现
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)
            throw new Error("Maximum lock count exceeded");
        setState(nextc);
        return true;
    }
    return false;
}

acquire 是 AQS 的 final 方法,逻辑很紧凑:

public final void acquire(int arg) {
    if (!tryAcquire(arg) &&                          // 先试一次
        acquireQueued(addWaiter(Node.EXCLUSIVE), arg)) // 失败则入队并阻塞等待
        selfInterrupt();                              // 等待期间被中断过,补一个中断标记
}

addWaiterenq 负责把当前线程包装成 Node 挂到队尾,用 CAS 保证并发安全:

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)) {     // 快速路径:一次 CAS 成功
            pred.next = node;
            return node;
        }
    }
    enq(node);                                   // 慢路径:自旋直到插入成功
    return node;
}

private Node enq(final Node node) {
    for (;;) {
        Node t = tail;
        if (t == null) {                          // 队列为空,先初始化头节点
            if (compareAndSetHead(new Node()))
                tail = head;
        } else {
            node.prev = t;
            if (compareAndSetTail(t, node)) {
                t.next = node;
                return t;
            }
        }
    }
}

入队之后的 acquireQueued 是整个机制的核心,一个无限循环:

final boolean acquireQueued(final Node node, int arg) {
    boolean failed = true;
    try {
        boolean interrupted = false;
        for (;;) {
            final Node p = node.predecessor();
            // 只有前驱是头节点才有资格尝试获取,保证 FIFO
            if (p == head && tryAcquire(arg)) {
                setHead(node);
                p.next = null;               // 断开引用,帮助 GC
                failed = false;
                return interrupted;
            }
            if (shouldParkAfterFailedAcquire(p, node) &&
                parkAndCheckInterrupt())
                interrupted = true;
        }
    } finally {
        if (failed)
            cancelAcquire(node);
    }
}

shouldParkAfterFailedAcquire 干的事是"确认前驱节点会在释放锁时通知我"。它检查前驱的 waitStatus,如果是 SIGNAL 就放心 park,如果是 CANCELLED 就往前跳过,如果是 0 就先 CAS 改成 SIGNAL

private static boolean shouldParkAfterFailedAcquire(Node pred, Node node) {
    int ws = pred.waitStatus;
    if (ws == Node.SIGNAL)
        return true;                          // 前驱说了会叫我,可以睡
    if (ws > 0) {
        do {                                  // 前驱放弃了,往前找到没放弃的
            node.prev = pred = pred.prev;
        } while (pred.waitStatus > 0);
        pred.next = node;
    } else {
        compareAndSetWaitStatus(pred, ws, Node.SIGNAL);   // 让前驱记得叫我
    }
    return false;
}

private final boolean parkAndCheckInterrupt() {
    LockSupport.park(this);                   // 真正挂起,底层是 Unsafe.park
    return Thread.interrupted();              // 注意这里会清除中断标记
}

公平和非公平的差别只有两行

公平版本的 lock() 直接走 acquire,一次都不抢:

// FairSync
static final class FairSync extends Sync {
    final void lock() {
        acquire(1);                           // 没有第一次抢
    }

    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;
    }
}

hasQueuedPredecessors() 判断"队列里是否有人比我等得更久":

public final boolean hasQueuedPredecessors() {
    Node t = tail;
    Node h = head;
    Node s;
    return h != t &&
        ((s = h.next) == null || s.thread != Thread.currentThread());
}

所以非公平锁一个线程有两次插队机会

  1. lock() 方法里的一次 CAS,完全不看队列
  2. 进入 acquiretryAcquire 里的又一次 CAS,还是不看队列

只有当这两次都失败,才老老实实排队。而在排队期间,新来的线程依然可以在它之前抢走锁。这就是我们那个 export-thread-7 被饿死的原因:它运气不好,每次轮到它之前都有新线程插队成功。

为什么非公平锁反而更快

既然会饥饿,为什么默认是它?因为吞吐差距很大。我写了个对比测试:20 个线程,每个线程循环 100 万次 lock() + 一次 i++ + unlock()

锁类型总耗时吞吐单线程获取次数(最大/最小)
非公平4.2 秒4.76 M ops/s1,284,391 / 47,206
公平21.7 秒0.92 M ops/s1,001,847 / 998,203

非公平的吞吐是公平的 5.2 倍,代价是各线程获取次数极不均匀,最少的那个只抢到 4.7 万次,最多的抢到 128 万次。

差距来自 park/unpark 的开销。实测一次 LockSupport.park + unpark 大约 1.2 微秒,涉及线程挂起、上下文切换、重新调度后缓存失效;而一次成功的 CAS 加锁只要 20 纳秒左右,差了 60 倍。

非公平锁的收益在于:刚释放锁的线程很可能还在 CPU 上跑(缓存是热的),让它立刻再拿到锁,比唤醒一个睡眠中的线程要快得多。而且在高并发下,"释放锁 → 唤醒下一个 → 下一个再抢"这个链条里,中间的空档期被插队者填满了,CPU 利用率更高。

释放锁:unparkSuccessor 为什么要从后往前找

public final boolean release(int arg) {
    if (tryRelease(arg)) {
        Node h = head;
        if (h != null && h.waitStatus != 0)
            unparkSuccessor(h);
        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);                           // 只有当前线程能调,不需要 CAS
    return free;
}

unparkSuccessor 里有个容易看不懂的地方:它从 tail 往前遍历找下一个要唤醒的节点,而不是顺着 next 走:

private void unparkSuccessor(Node node) {
    int ws = node.waitStatus;
    if (ws < 0)
        compareAndSetWaitStatus(node, ws, 0);

    Node s = node.next;
    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);
}

原因是 next 指针在节点被取消时是不可靠的。addWaiter 里先执行 node.prev = pred 再 CAS 设置 tail,最后才 pred.next = node,这三步之间如果有节点被取消,next 链可能断掉;而 prev 链因为始终是"先连前驱再改 tail"的顺序,相对稳定。Doug Lea 在注释里专门写了这一点。

我们最后的选择

回到导出那个场景。它的特点是临界区比较长:写一行 Excel 涉及 14 个单元格的创建,实测持锁时间约 0.8 毫秒。这种情况下,非公平锁的插队收益被摊薄了(因为线程大部分时间在干活,不是在抢锁),而饥饿的代价被放大。

改成公平锁之后:

指标非公平公平
20 万行导出总耗时3 分 12 秒3 分 40 秒
单线程最长累计等待41.3 秒6.8 秒
各线程等待时间标准差13.7 秒0.4 秒

总耗时慢了 15%,但最慢的那个线程少等了 34 秒,各线程的负载也均匀了。我们接受这个交换。

判断标准其实很简单:临界区短(纳秒到微秒级)用非公平,临界区长(毫秒级以上)且线程数多用公平。另外公平锁不支持"闯入",所以它在负载低的时候也没法优化,吞吐一定更差。

写在后面

现在回头看,《AQS 源码初探:ReentrantLock 是怎么实现的》本身不算多难,难的是线上真出问题那十分钟里的判断。经验都是这么来的。

参考