CountDownLatch
CountDownLatch 源码分析
一次性的”倒计时门闩”:初始化计数 N,各线程 countDown() 减一,await() 的线程阻塞到计数归零后全部放行。
核心思路是把 AQS 的 state 直接当计数用——归零即开门,且是永久开门,不能重置(要循环用的场景是 CyclicBarrier 的地盘)。
整个类不到 60 行有效代码,是 AQS 共享模式最干净的教科书用法。
代码块收起展开
// 基于 JDK 17 (本地 ms-17.0.18), java.util.concurrent.CountDownLatch
public class CountDownLatch {
// 内部同步器。CountDownLatch 自己不写任何并发逻辑,全部委托给 AQS
private static final class Sync extends AbstractQueuedSynchronizer {
private static final long serialVersionUID = 4982264981922014374L;
Sync(int count) {
setState(count); // state 被赋予"剩余计数"的含义,构造时一次写入,之后只减不增
}
int getCount() {
return getState();
}
// 共享获取:语义是"问门开没开",不消耗任何资源,所以只读不写、没有 CAS。
// 返回 1 = 门开(state==0)直接过;返回 -1 = 没开,回 AQS 入队挂起
protected int tryAcquireShared(int acquires) {
return (getState() == 0) ? 1 : -1;
}
// 共享释放:CAS 自旋减一。整个类唯一有竞争的地方就在这
protected boolean tryReleaseShared(int releases) {
// Decrement count; signal when transition to zero
for (;;) {
int c = getState();
if (c == 0)
return false; // 已归零就直接失败:多余的 countDown 无害,state 永不为负,也不会重复触发唤醒
int nextc = c - 1;
if (compareAndSetState(c, nextc))
return nextc == 0; // 只有"减到 0 的那一次"返回 true → AQS 才去 signalNext。N 次 countDown 只有一次触发唤醒
}
}
}
private final Sync sync;
public CountDownLatch(int count) {
if (count < 0) throw new IllegalArgumentException("count < 0");
this.sync = new Sync(count); // count==0 合法:等价于一个天生开着的门,await 直接过
}
// 阻塞等待归零,可中断。用的是共享+可中断版获取
public void await() throws InterruptedException {
sync.acquireSharedInterruptibly(1);
}
// 带超时版:超时返回 false 而不是永远卡死,生产代码里几乎总该用这个
public boolean await(long timeout, TimeUnit unit)
throws InterruptedException {
return sync.tryAcquireSharedNanos(1, unit.toNanos(timeout));
}
// 减一,归零时放行所有 await 线程。注意它自己从不阻塞
public void countDown() {
sync.releaseShared(1);
}
public long getCount() {
return sync.getCount();
}
// ... 省略 toString
}await/countDown 各自只有一行委托,真正的排队与唤醒在 AQS 里。下面是共享模式在 JDK 17 里的骨架(AQS 在 JDK 17 被 Doug Lea 重写过,入队+阻塞全部收进一个 acquire 方法):
代码块收起展开
// 基于 JDK 17 (本地 ms-17.0.18), java.util.concurrent.locks.AbstractQueuedSynchronizer
public abstract class AbstractQueuedSynchronizer
extends AbstractOwnableSynchronizer
implements java.io.Serializable {
// ...
// await() 的入口。注意中断检查在 tryAcquireShared 之前:进门先看中断标记,被中断过就不排队直接抛
public final void acquireSharedInterruptibly(int arg)
throws InterruptedException {
if (Thread.interrupted() ||
(tryAcquireShared(arg) < 0 &&
acquire(null, arg, true, true, false, 0L) < 0)) // shared=true, interruptible=true;返回负数=排队期间被中断
throw new InterruptedException();
}
// countDown() 的入口
public final boolean releaseShared(int arg) {
if (tryReleaseShared(arg)) { // 只有归零那一次为 true
signalNext(head); // 注意:只 unpark 队头后继"一个"线程,不是全部
return true;
}
return false;
}
// 唤醒 h 的后继。先清 WAITING 状态再 unpark:被唤醒线程靠状态位区分"真信号"和虚假唤醒,防 park 竞态
private static void signalNext(Node h) {
Node s;
if (h != null && (s = h.next) != null && s.status != 0) {
s.getAndUnsetStatus(WAITING);
LockSupport.unpark(s.waiter);
}
}
// 和 signalNext 唯一区别:只对 SharedNode 生效,独占节点不传播
private static void signalNextIfShared(Node h) {
Node s;
if (h != null && (s = h.next) != null &&
(s instanceof SharedNode) && s.status != 0) {
s.getAndUnsetStatus(WAITING);
LockSupport.unpark(s.waiter);
}
}
// 入队+阻塞的统一实现。只看共享模式获取成功后的那段:
final int acquire(Node node, int arg, boolean shared,
boolean interruptible, boolean timed, long time) {
// ...
if (acquired) {
if (first) {
node.prev = null;
head = node; // 自己成为新 head
pred.next = null;
node.waiter = null;
if (shared)
signalNextIfShared(node); // 关键:共享模式下,我过门后顺手叫醒下一个 → 唤醒沿队列链式传播
if (interrupted)
current.interrupt();
}
return 1;
}
// ...
}
}原理串讲
以 new CountDownLatch(3)、一个主线程 await()、三个工作线程 countDown() 为例走一遍完整链路。
构造时 Sync(3) 调 setState(3),state 就是剩余计数。主线程调 await() → acquireSharedInterruptibly(1):先查中断标记,然后 tryAcquireShared(1) 读到 state==3 != 0 返回 -1,于是进 acquire(null, 1, true, true, false, 0L)——在等待队列里挂一个 SharedNode,把自己 LockSupport.park 住。
多个 await 线程就在队列里排成一串共享节点。
工作线程调 countDown() → releaseShared(1) → tryReleaseShared(1) 在 CAS 自旋里把 state 减一。
前两次 countDown 把 state 从 3 减到 1,nextc != 0 返回 false,releaseShared 直接返回——门没开,什么唤醒都不发生,这也是 countDown 永不阻塞的原因。
第三次减到 0 返回 true,releaseShared 调 signalNext(head),unpark 队头后继那一个线程。
被唤醒的主线程从 acquire 的 park 处醒来,重试 tryAcquireShared,这次 state<mark>0 返回 1,获取成功:把自己设为新 head,然后因为 shared</mark>true 调 signalNextIfShared(node) 唤醒下一个共享节点。
下一个醒来重复同样的动作。唤醒像多米诺骨牌一样沿队列传播,直到所有 await 线程全部放行。
之后再来的 await 在快速路径 tryAcquireShared 就直接返回 1,连队列都不进。
为什么归零时只 signalNext 唤醒一个,而不是遍历队列全部唤醒? 因为唤醒的责任被摊派了:每个被唤醒的线程在成为 head 之后自己去叫醒下一个(signalNextIfShared)。
执行 countDown 的线程 O(1) 就能返回,不用背着”遍历唤醒整条队列”的开销;同时链式传播天然按 FIFO 顺序展开,逻辑上等价于广播,成本上是流水线。
为什么 tryAcquireShared 只读不写、连 CAS 都没有? 独占锁的”获取”要占坑(改 state 记录持有者),而 latch 的”获取”只是探询门开没开,不消耗任何资源、不修改任何状态——门一旦开了所有人都能过,读一次 volatile 的 state 就够。
这就是共享模式和独占模式的本质区别,也是它比 ReentrantLock 的 tryAcquire 简单得多的原因。
为什么 c == 0 时 tryReleaseShared 返回 false 而不是抛异常? 幂等保护。
业务代码里 finally 里多调一次 countDown、或异常路径重复计数是常态,Doug Lea 选择让多余的 countDown 静默无害:state 不会减成负数,也不会重复触发 signalNext。
代价是”计数错了少减一次”这种 bug 不会有任何报错,只会表现为 await 永远醒不来——所以生产上优先用带超时的 await(timeout, unit) 兜底。
最后一个常被问到的点:countDown() 之前的写操作 happen-before 另一个线程 await() 返回之后的读。
载体就是 volatile 的 state——compareAndSetState 的释放语义和 getState 的获取语义搭出了这条内存可见性通道,所以工作线程准备好的数据,主线程 await() 醒来后直接读是安全的,不需要额外同步。
设计取舍
- 一次性 vs 可重置:state 单向递减、归零永久开门,换来实现极简。若允许 reset,会和正在链式传播的唤醒产生竞态,CyclicBarrier 靠 lock + condition + 换代(generation)才做到可重复用,复杂度翻几倍。
- countDown 与 await 完全解耦:countDown 的线程从不等待、从不阻塞,减完就走。这让 latch 既能”1 等 N”(汇总)也能”N 等 1”(发令枪,count=1)。
- 多余的 countDown 静默吞掉:换来幂等安全,代价是计数配错时无声死等——用带超时的 await 兜底。
- 常见误区:
await可被中断抛InterruptedException,吞掉不处理会让线程带着假状态继续跑;getCount()只适合调试,拿它做业务判断必有竞态。 - AQS 共享模式的完整骨架见 AbstractQueuedSynchronizer。