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) | 超时和正常放行无差别,无法感知异步是否真完成 |
| 2 | InterruptedException 被吞 | await 响应中断抛 IE(§2.3 parkAndCheckInterrupt) | 只打日志,且未恢复中断标志(Thread.currentThread().interrupt()),常见反模式 |
| 3 | run() 无 try-catch 放大效应 | 若 execute() 抛异常,countDown() 不执行 → state 永远 >0 | await 只能靠 1000ms 超时返回(§2.4 超时分支),主线程白白阻塞 1s |
| 4 | countDown 在 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需重新notify | count 不可重置(一次性闸门) |
| 底层 | JVM monitor wait | AQS CLH 队列 + LockSupport.park |
九、源码定位(本地查看完整源码)
| 类 | JDK 源码路径 |
|---|---|
CountDownLatch | java.util.concurrent.CountDownLatch |
Sync(内部类) | 同上,CountDownLatch 内部 |
AbstractQueuedSynchronizer | java.util.concurrent.locks.AbstractQueuedSynchronizer |
LockSupport | java.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);绿色 = 唤醒点(发许可)。