06 - CountDownLatch 源码解析与 AQS 机制

本篇解析启动框架 AppStartTaskDispatcher 用到的 CountDownLatch.await(timeout) 的 JDK 源码,从 CountDownLatch 一路下沉到 AQS、LockSupport,直至 native。 关联:01-启动框架与调度机制 § CountDownLatch 同步模型


〇、命名澄清

CountDownLatch 没有 wait() 方法——wait()Object 的方法(配合 synchronized/notify)。CountDownLatch 的等待 API 是 await() / await(long, TimeUnit)

项目实际用的是带超时版(AppStartTaskDispatcher.java:75):

mCountDownLatch.await(mAllTaskWaitTimeOut, TimeUnit.MILLISECONDS);  // 超时 1000ms

源码版本:本文基于 JDK 8(Android Core / OpenJDK 核心一致)。JDK 9+ 把 tryAcquireSharedNanos 的内联实现提取成了私有方法 doAcquireSharedNanos,逻辑相同。


一、核心数据结构

CountDownLatch 本身很薄,核心逻辑全在内部类 Sync(继承 AQS)。用 AQS 的 state 表示剩余计数

public class CountDownLatch {
    private final Sync sync;
 
    // 构造:把 count 装进 AQS 的 state
    public CountDownLatch(int count) {
        if (count < 0) throw new IllegalArgumentException("count < 0");
        sync = new Sync(count);
    }
 
    private static final class Sync extends AbstractQueuedSynchronizer {
        Sync(int count) { setState(count); }              // state = count
        int getCount() { return getState(); }
 
        // "是否放行":state==0 才放行(返回≥0),否则阻塞(返回-1)
        protected int tryAcquireShared(int acquires) {
            return (getState() == 0) ? 1 : -1;
        }
 
        // "递减计数":CAS 把 state-1,减到 0 时返回 true 触发唤醒
        protected boolean tryReleaseShared(int releases) {
            for (;;) {
                int c = getState();
                if (c == 0) return false;                 // 已经是0,不再减
                int nextc = c - 1;
                if (compareAndSetState(c, nextc))         // CAS 自旋
                    return nextc == 0;                    // 减到0才返回true
            }
        }
    }
}

关键设计:

  • tryAcquireShared 只看 state==0 —— 不消费 state,只是”探测闸门是否打开”。这就是为什么多个线程可同时 await,且 countDown 到 0 后唤醒所有等待者。
  • tryReleaseShared 用 CAS 自旋递减,保证多线程并发 countDown 安全。
  • CountDownLatch 的 state 不可重置(构造时一次设定),区别于 CyclicBarrier 可循环使用。

二、await() 源码链路

2.1 两个 await 入口

// 无限等待(可中断)
public void await() throws InterruptedException {
    sync.acquireSharedInterruptibly(1);
}
 
// 带超时等待(项目用的这个) —— 返回 true=计数到0放行;false=超时
public boolean await(long timeout, TimeUnit unit) throws InterruptedException {
    return sync.tryAcquireSharedNanos(1, unit.toNanos(timeout));
}

2.2 AQS.acquireSharedInterruptibly(无限版)

// java.util.concurrent.locks.AbstractQueuedSynchronizer
public final void acquireSharedInterruptibly(int arg) throws InterruptedException {
    if (Thread.interrupted())                              // 1.响应中断
        throw new InterruptedException();
    if (tryAcquireShared(arg) < 0)                         // 2.CountDownLatch 实现:state!=0 返回-1
        doAcquireSharedInterruptibly(arg);                 // 3.入队阻塞
}

2.3 AQS.doAcquireSharedInterruptibly(真正的阻塞逻辑)

private void doAcquireSharedInterruptibly(int arg) throws InterruptedException {
    final Node node = addWaiter(Node.SHARED);              // 1.包装成 SHARED 节点加入 CLH 队列尾部
    boolean failed = true;
    try {
        for (;;) {                                         // 2.自旋
            final Node p = node.predecessor();             //   取前驱
            if (p == head) {                               //   前驱是头,说明轮到自己了
                int r = tryAcquireShared(arg);             //   再探一次 state==0?
                if (r >= 0) {                              //   放行
                    setHeadAndPropagate(node, r);          //   ★共享传播:唤醒后续共享节点
                    p.next = null;                         //   help GC
                    failed = false;
                    return;
                }
            }
            // 3.没轮到/没放行 → 判断是否应阻塞
            if (shouldParkAfterFailedAcquire(p, node) &&   //   把前驱的 waitStatus 置为 SIGNAL
                parkAndCheckInterrupt())                   //   ★LockSupport.park 阻塞当前线程
                throw new InterruptedException();          //   阻塞期间被中断则抛异常
        }
    } finally {
        if (failed) cancelAcquire(node);                   // 4.异常时取消节点
    }
}
 
private final boolean parkAndCheckInterrupt() {
    LockSupport.park(this);                                // ★真正挂起线程的地方
    return Thread.interrupted();
}

2.4 带超时版 doAcquireSharedNanos

带超时版多了 deadline 计算parkNanos,以及一个优化:nanosTimeout > SPIN_FOR_TIMEOUT_THRESHOLD(1000ns) 才 park,否则自旋(避免极短超时的 park 唤醒开销):

private boolean doAcquireSharedNanos(int arg, long nanosTimeout) throws InterruptedException {
    if (nanosTimeout <= 0L) return false;                  // 超时直接返回false
    final Node node = addWaiter(Node.SHARED);
    boolean failed = true;
    try {
        final long deadline = System.nanoTime() + nanosTimeout;   // ★计算截止时刻
        for (;;) {
            final Node p = node.predecessor();
            if (p == head) {
                int r = tryAcquireShared(arg);
                if (r >= 0) {
                    setHeadAndPropagate(node, r);
                    p.next = null; failed = false;
                    return true;                           // ★放行返回true
                }
            }
            nanosTimeout = deadline - System.nanoTime();   // 剩余时间
            if (nanosTimeout <= 0L) return false;          // ★超时返回false
            if (shouldParkAfterFailedAcquire(p, node) && nanosTimeout > SPIN_FOR_TIMEOUT_THRESHOLD)
                LockSupport.parkNanos(this, nanosTimeout); // ★限时阻塞
            if (Thread.interrupted()) throw new InterruptedException();
        }
    } finally {
        if (failed) cancelAcquire(node);
    }
}

三、countDown() 唤醒链路

public void countDown() {
    sync.releaseShared(1);
}

3.1 AQS.releaseShared

public final boolean releaseShared(int arg) {
    if (tryReleaseShared(arg)) {        // 1.CountDownLatch 实现:CAS递减,到0返回true
        doReleaseShared();              // 2.★唤醒队列中等待的线程
        return true;
    }
    return false;
}

3.2 AQS.doReleaseShared(传播唤醒)

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;
                unparkSuccessor(h);                     // ★LockSupport.unpark 唤醒后继
            } else if (ws == 0 &&
                       !compareAndSetWaitStatus(h, 0, Node.PROPAGATE))
                continue;
        }
        if (h == head) break;
    }
}

unpark 唤醒的等待线程从 parkAndCheckInterrupt 返回,回到 for(;;) 自旋顶部,重试 tryAcquireShared,此时 state==0 返回 1 ≥ 0,setHeadAndPropagate 放行并继续向后传播唤醒(共享模式的核心:一次 countDown(到0) 唤醒所有 await 者)。


四、底层:LockSupport.park / unpark

// java.util.concurrent.locks.LockSupport
public static void park(Object blocker) {
    Thread t = Thread.currentThread();
    setBlocker(t, blocker);
    U.park(false, 0L);                  // → Unsafe#park (native)
    setBlocker(t, null);
}
public static void unpark(Thread thread) {
    if (thread != null)
        U.unpark(thread);               // → Unsafe#unpark (native)
}
  • park 挂起当前线程(许可语义:若已有许可则立即返回并消耗,否则阻塞)
  • unpark(thread) 给目标线程发一张许可(先发后用也有效,这是与 Object.notify 的关键区别——notify 必须在 wait 之后调用才有效)
  • 底层是 JVM native,对应操作系统的 futex(Linux)/ park primitive

五、整体交互时序

sequenceDiagram
    participant Main as 主线程(await)
    participant Sync as CountDownLatch.Sync
    participant AQS as AQS(CLB队列)
    participant Async as 异步线程(countDown)
    participant LS as LockSupport

    Main->>Sync: await(1000, MS)
    Sync->>AQS: tryAcquireSharedNanos
    AQS->>Sync: tryAcquireShared: state=1 返回 -1
    Note over AQS: state 不等于 0,放行失败
    AQS->>AQS: addWaiter(SHARED) 入队尾
    AQS->>LS: parkNanos(1000ms) 阻塞主线程
    Note over Main: 主线程挂起(释放CPU)

    Async->>Sync: countDown()
    Sync->>Sync: tryReleaseShared: CAS state 1到0 返回true
    Sync->>AQS: doReleaseShared
    AQS->>LS: unpark(主线程) 发许可
    LS-->>Main: 唤醒(从 parkNanos 返回)
    Main->>AQS: 自旋重试 tryAcquireShared
    AQS->>Sync: state==0 返回 1
    Note over AQS: r 大于等于 0,放行
    AQS->>AQS: setHeadAndPropagate 传播唤醒
    AQS-->>Main: await 返回 true

六、AQS CLH 队列状态流转

stateDiagram-v2
    [*] --> New: new Node(SHARED)
    New --> Enqueued: addWaiter 加入队尾
    Enqueued --> Spinning: for(;;) 自旋
    Spinning --> TryAcquire: 前驱==head
    TryAcquire --> Passed: state==0 返回≥0
    TryAcquire --> CheckPark: state!=0 返回-1
    CheckPark --> Parked: shouldParkAfterFailedAcquire
    Parked --> Spinning: 被 unpark 唤醒
    Parked --> Interrupted: 阻塞期被中断
    Passed --> Propagate: setHeadAndPropagate
    Propagate --> [*]: return (放行)
    Interrupted --> [*]: throw InterruptedException
    note right of Parked
        LockSupport.parkNanos
        真正释放 CPU 的地方
    end note

七、结合项目的源码级注意点

回到 AppStartTaskDispatcher.java:65-79:

try {
    if (mAllTaskWaitTimeOut == 0) mAllTaskWaitTimeOut = WAITING_TIME;
    mCountDownLatch.await(mAllTaskWaitTimeOut, TimeUnit.MILLISECONDS);
} catch (InterruptedException e) {
    MLog.e(e, "AppStartTaskDispatcher await()");
}
#问题源码依据影响
1返回值被忽略await(timeout) 返回 boolean(true=放行 / false=超时,见 §2.4)超时和正常放行无差别,无法感知异步是否真完成
2InterruptedException 被吞await 响应中断抛 IE(§2.3 parkAndCheckInterrupt)只打日志,且未恢复中断标志(Thread.currentThread().interrupt()),常见反模式
3run() 无 try-catch 放大效应execute() 抛异常,countDown() 不执行 → state 永远 >0await 只能靠 1000ms 超时返回(§2.4 超时分支),主线程白白阻塞 1s
4countDown 在 markFinish 里AbsStartTask.run()execute()markFinish → countDown顺序保证:业务异常会阻断 countDown(见 #3)

八、Object.wait() vs CountDownLatch.await()

维度Object.wait()CountDownLatch.await()
依赖必须持有该对象的 synchronized基于 AQS,无需 synchronized
唤醒notify/notifyAll,且唤醒时必须重新竞争锁countDown 到 0,基于 park/unpark 许可
条件需自己写 while 循环防虚假唤醒内部封装好(state==0 判定)
复用每次wait需重新notifycount 不可重置(一次性闸门)
底层JVM monitor waitAQS CLH 队列 + LockSupport.park

九、源码定位(本地查看完整源码)

JDK 源码路径
CountDownLatchjava.util.concurrent.CountDownLatch
Sync(内部类)同上,CountDownLatch 内部
AbstractQueuedSynchronizerjava.util.concurrent.locks.AbstractQueuedSynchronizer
LockSupportjava.util.concurrent.locks.LockSupport
  • JDK 源码包:$JAVA_HOME/lib/src.zip(解压后 java/util/concurrent/)
  • Android Studio:External Libraries → <JDK>/java.util.concurrent/CountDownLatch.class 反编译,或 Attach Sources
  • 在线:OpenJDK 仓库 jdk/src/share/classes/java/util/concurrent/CountDownLatch.java

十、一图总结调用链

flowchart TD
    A["CountDownLatch.await(timeout)"] --> B["sync.tryAcquireSharedNanos"]
    A2["CountDownLatch.await"] --> B2["sync.acquireSharedInterruptibly"]
    B --> C["AQS.doAcquireSharedNanos"]
    B2 --> C2["AQS.doAcquireSharedInterruptibly"]
    C --> D["tryAcquireShared: state==0?"]
    C2 --> D
    D -->|"state!=0 返回-1"| E["addWaiter 入CLH队列"]
    E --> F["LockSupport.parkNanos / park"]
    F -->|"unpark 唤醒"| D

    G["CountDownLatch.countDown"] --> H["sync.releaseShared"]
    H --> I["tryReleaseShared: CAS state-1"]
    I -->|"减到0"| J["AQS.doReleaseShared"]
    J --> K["unparkSuccessor"]
    K --> L["LockSupport.unpark(等待线程)"]
    L -.->|"许可"| F

    style F fill:#ffe1e1,stroke:#cc0000
    style L fill:#e1f5e1,stroke:#008800

红色 = 阻塞点(释放 CPU);绿色 = 唤醒点(发许可)。