AQS 共享模式只服务两件:CountDownLatch 和 Semaphore
一个订单导出接口要并发调三个下游,等齐了再拼装。需求一句话,可选型一摊开是一整族:CountDownLatch、Semaphore、CyclicBarrier、Phaser、Exchanger、ReentrantReadWriteLock、StampedLock。面试里关于它们最高频的一句错答是"这几个都基于 AQS 共享模式"——这句话只在其中两个身上成立,其余五个连 AQS 都没继承。
这篇文章用一个能跑的导出场景把这七个类串起来,每个结论都配一段可以自己复现的输出。源码引用基于实测 JDK 26,碰到版本差异的地方会单独标出来。
先把场景立起来:一个导出接口,三个下游
导出要拿到用户、商品、库存三份数据,三份齐了才能拼装,任一下游超时就整体失败。最直接的写法是 CountDownLatch(3),三个任务各自干完就 countDown,主线程 await。
先看三个共用的辅助方法,后面三段都依赖它们:
static String load(String name, long ms) {
try {
Thread.sleep(ms);
} catch (InterruptedException e) {
return name + ":interrupted";
}
return name + ":ok";
}
static String loadFail(String name, long ms) {
try {
Thread.sleep(ms);
} catch (InterruptedException e) {
return name + ":interrupted";
}
throw new IllegalStateException(name + " down");
}
static String assemble(String u, String i, String s) {
return "assemble(" + u + ", " + i + ", " + s + ")";
}
第一版是顺利路径,三个下游分别耗时 20ms、40ms、60ms,超时给足 2 秒:
static void happyPath() throws Exception {
CountDownLatch latch = new CountDownLatch(3);
AtomicReference<String> u = new AtomicReference<>();
AtomicReference<String> i = new AtomicReference<>();
AtomicReference<String> s = new AtomicReference<>();
ExecutorService pool = Executors.newFixedThreadPool(3);
pool.submit(() -> { try { u.set(load("user", 20)); } finally { latch.countDown(); } });
pool.submit(() -> { try { i.set(load("item", 40)); } finally { latch.countDown(); } });
pool.submit(() -> { try { s.set(load("stock", 60)); } finally { latch.countDown(); } });
boolean ok = latch.await(2, TimeUnit.SECONDS);
System.out.println("[1] awaitOk=" + ok + " -> " + assemble(u.get(), i.get(), s.get()));
pool.shutdown();
pool.awaitTermination(2, TimeUnit.SECONDS);
}
两个细节值得停一下。第一,countDown() 写在 finally 里,不是可有可无的谨慎——第三版会证明漏掉它的后果。第二,load 把数据塞进 AtomicReference 再让主线程读,因为 latch.await() 返回只保证"计数到零",它与三个工作线程的写之间有 happens-before 关系(AQS 的 state 是 volatile,countDown 里用 CAS 改它),所以这里不靠额外的锁也能安全读到。用普通数组或 HashMap 装结果就是在赌运气。
第二版把库存改成 500ms,主线程只等 100ms:
static void timeoutPath() throws Exception {
CountDownLatch latch = new CountDownLatch(3);
AtomicReference<String> s = new AtomicReference<>();
ExecutorService pool = Executors.newFixedThreadPool(3);
pool.submit(() -> { try { load("user", 20); } finally { latch.countDown(); } });
pool.submit(() -> { try { load("item", 40); } finally { latch.countDown(); } });
pool.submit(() -> { try { s.set(load("stock", 500)); } finally { latch.countDown(); } });
boolean ok = latch.await(100, TimeUnit.MILLISECONDS);
System.out.println("[2] awaitOk=" + ok + " stockBeforeTimeout=" + s.get());
boolean ok2 = latch.await(2, TimeUnit.SECONDS);
System.out.println("[2] awaitOkSecondTry=" + ok2 + " stockAfter=" + s.get());
pool.shutdown();
pool.awaitTermination(2, TimeUnit.SECONDS);
}
第三版让库存接口直接抛异常,而任务里的 catch 忘了 countDown:
static void stuckPath() throws Exception {
CountDownLatch latch = new CountDownLatch(3);
ExecutorService pool = Executors.newFixedThreadPool(3);
pool.submit(() -> { try { load("user", 20); } finally { latch.countDown(); } });
pool.submit(() -> { try { load("item", 40); } finally { latch.countDown(); } });
pool.submit(() -> {
try {
loadFail("stock", 20); // 抛异常
} catch (IllegalStateException e) {
// 这里忘了写 finally { latch.countDown(); }
}
});
boolean ok = latch.await(300, TimeUnit.MILLISECONDS);
System.out.println("[3] awaitOk=" + ok + " latchCountStuckAt=" + latch.getCount()
+ " (the third task threw before it could count down)");
pool.shutdown();
pool.awaitTermination(2, TimeUnit.SECONDS);
}
main 里按顺序调这三段,下面是完整输出:
[1] awaitOk=true -> assemble(user:ok, item:ok, stock:ok)
[2] awaitOk=false stockBeforeTimeout=null
[2] awaitOkSecondTry=true stockAfter=stock:ok
[3] awaitOk=false latchCountStuckAt=1 (the third task threw before it could count down)
四行输出正好把 CountDownLatch 的三个关键行为摆出来了。await(timeout) 返回 boolean,超时那一刻 stockBeforeTimeout=null,说明子任务确实还没写完;但计数不会因此归零,再等一次就拿到了 stock:ok——超时只影响这一次 await,不影响计数器,这是它和 CyclicBarrier"破损"语义的根本分歧。最后一行更狠:计数永久停在 1,主线程每次 await 都只能等到超时。生产里这类事故的现场就是"接口不报错,只是每次都慢到超时"。
一张图看清谁站在 AQS 上
先把家族的骨架画出来,再逐个查。左边这一支继承 AbstractQueuedSynchronizer,右边那一档压根不走 AQS:
图里最上面那个 AbstractOwnableSynchronizer 容易被忽略:它只记一件事——独占锁当前握在哪个线程手里,是 JDK 6 才抽出来的基类。AQS 和 AQLS 都继承它,所以两者都有 getExclusiveOwnerThread()。
"基于 AQS"这句话需要拆成两层看。AbstractQueuedSynchronizer 用 int state,AbstractQueuedLongSynchronizer 用 long state,两个是平级的兄弟。落在 AQS 上的只有 CountDownLatch、Semaphore(共享模式)和 ReentrantLock(独占模式);落在 AQLS 上的是 JDK 26 之后的 ReentrantReadWriteLock.Sync。而 CyclicBarrier、Phaser、StampedLock、Exchanger 各有各的实现,一个都不沾。
光看图不够,跑一段代码把每个类对外暴露的"state 派生值"列出来对比:
public class StateView {
public static void main(String[] args) throws Exception {
CountDownLatch latch = new CountDownLatch(3);
latch.countDown();
System.out.println("[CountDownLatch] getCount=" + latch.getCount() + " (state itself)");
Semaphore sem = new Semaphore(5);
sem.acquire(2);
System.out.println("[Semaphore] availablePermits=" + sem.availablePermits()
+ " queueLength=" + sem.getQueueLength() + " (state itself)");
ReentrantReadWriteLock rw = new ReentrantReadWriteLock();
rw.writeLock().lock();
System.out.println("[ReentrantReadWriteLock] writeHoldCount=" + rw.getWriteHoldCount()
+ " readLockCount=" + rw.getReadLockCount() + " isWriteLocked=" + rw.isWriteLocked());
rw.readLock().lock(); // 持写锁时拿读锁(降级)
rw.writeLock().unlock();
System.out.println("[ReentrantReadWriteLock] afterDowngrade readLockCount="
+ rw.getReadLockCount() + " readHoldCount=" + rw.getReadHoldCount()
+ " isWriteLocked=" + rw.isWriteLocked());
rw.readLock().unlock();
StampedLock stamped = new StampedLock();
long ws = stamped.writeLock();
System.out.println("[StampedLock] writeStampLow8=0x" + Long.toHexString(ws & 0xFF).toUpperCase()
+ " isWriteLocked=" + stamped.isWriteLocked()
+ " readLockCount=" + stamped.getReadLockCount());
stamped.unlockWrite(ws);
Phaser phaser = new Phaser(2);
System.out.println("[Phaser] registeredParties=" + phaser.getRegisteredParties()
+ " phase=" + phaser.getPhase() + " arrivedParties=" + phaser.getArrivedParties());
phaser.arriveAndDeregister();
System.out.println("[Phaser] afterDeregister registeredParties=" + phaser.getRegisteredParties()
+ " terminated=" + phaser.isTerminated());
CyclicBarrier barrier = new CyclicBarrier(3);
System.out.println("[CyclicBarrier] parties=" + barrier.getParties()
+ " numberWaiting=" + barrier.getNumberWaiting() + " isBroken=" + barrier.isBroken());
}
}
[CountDownLatch] getCount=2 (state itself)
[Semaphore] availablePermits=3 queueLength=0 (state itself)
[ReentrantReadWriteLock] writeHoldCount=1 readLockCount=0 isWriteLocked=true
[ReentrantReadWriteLock] afterDowngrade readLockCount=1 readHoldCount=1 isWriteLocked=false
[StampedLock] writeStampLow8=0x80 isWriteLocked=true readLockCount=0
[Phaser] registeredParties=2 phase=0 arrivedParties=0
[Phaser] afterDeregister registeredParties=1 terminated=false
[CyclicBarrier] parties=3 numberWaiting=0 isBroken=false
前两行都在说"这就是 state 本身":CountDownLatch.getCount() 直接返回 AQS 的 state,Semaphore.availablePermits() 也是。所以这两个类的对外语义和内部 state 是同一件事,读懂 state 就读懂了它们。
再往下就不一样了。ReentrantReadWriteLock 一次要回答三个问题(读计数、写重入次数、写锁是否被持有),一个 int 不够用,只能把 state 切成两半分别装。StampedLock 更直接——它不设 getter,你得自己 stamp & 0xFF 去解,因为 stamp 就是 state 的快照。Phaser 和 CyclicBarrier 的 getter 返回的是纯逻辑量(参与方、阶段号、等待人数),底下没有任何 AQS state 与之对应。
CountDownLatch:state 就是一个倒计数器
CountDownLatch 的全部秘密在内部类 Sync 里,JDK 26 的源码是这样的:
private static final class Sync extends AbstractQueuedSynchronizer {
Sync(int count) {
setState(count);
}
int getCount() {
return getState();
}
protected int tryAcquireShared(int acquires) {
return (getState() == 0) ? 1 : -1;
}
protected boolean tryReleaseShared(int releases) {
// Decrement count; signal when transition to zero
for (;;) {
int c = getState();
if (c == 0)
return false;
int nextc = c - 1;
if (compareAndSetState(c, nextc))
return nextc == 0;
}
}
}
三个方法对应三件事。构造时把 count 灌进 state。countDown() 走 releaseShared(1) → tryReleaseShared,CAS 把 state 减一;注意它的返回值是 nextc == 0——只有减到零那一次才返回 true,AQS 收到 true 才会去唤醒等待队列,前面的 2、1 都返回 false 什么也不做。await() 走 acquireSharedInterruptibly(1) → tryAcquireShared,state == 0 ? 1 : -1:state 不为零就返回负数,AQS 把当前线程挂进共享模式的队列。
为什么必须是共享模式而不是独占?因为 await() 可以被多个线程同时调用(主线程和定时任务都等同一批下游数据是常见写法)。state 归零那一次,共享队列上的所有等待节点要一起醒来;独占模式一次只放一个线程,第一个醒来的线程不释放就没完没了。共享模式的实现要点是节点唤醒后会继续向后传播唤醒动作,AQS 里用 PROPAGATE 这个中间状态来避免"释放了但没人醒"的丢失唤醒问题。
回到前面的第三版:latch.getCount() 卡在 1,就是 tryReleaseShared 从来没被调用过第三次,nextc == 0 这个条件永远等不到。CountDownLatch 没有超时兜底、没有重置、没有破损检测——计数到不了零就是永久等待,唯一的安全网是 await(timeout) 的布尔返回值,而它需要调用方自己检查。这就是为什么 countDown() 必须进 finally。
Semaphore:state 是许可池
Semaphore 和 CountDownLatch 是同一套骨架的两种用法。看下面这段,三个许可配八个任务:
static void permits() throws Exception {
Semaphore sem = new Semaphore(3);
AtomicInteger cur = new AtomicInteger();
AtomicInteger peak = new AtomicInteger();
ExecutorService pool = Executors.newFixedThreadPool(8);
CountDownLatch done = new CountDownLatch(8);
for (int k = 0; k < 8; k++) {
pool.submit(() -> {
try {
sem.acquire();
try {
int c = cur.incrementAndGet();
peak.accumulateAndGet(c, Math::max);
sleep(40);
} finally {
cur.decrementAndGet();
sem.release();
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
} finally {
done.countDown();
}
return null;
});
}
done.await(10, TimeUnit.SECONDS);
pool.shutdown();
pool.awaitTermination(5, TimeUnit.SECONDS);
System.out.println("[1] peakConcurrent=" + peak.get() + " availablePermits=" + sem.availablePermits()
+ " queueLength=" + sem.getQueueLength());
}
这里有个容易忽略的点:peak 统计的是"同时持有许可的任务数",用 AtomicInteger 自增自减加 Math::max 累积。如果直接用 System.out.println 在任务里打印,八个线程的输出会交织成乱码——并发演示想让输出稳定,就不能让工作线程自己打印,要攒到计数器里最后一次性输出。这是所有并发 Demo 能不能复现的前提。
再看两个变体:许可给 1 就是互斥锁,tryAcquire 拿不到直接返回不排队。
static void asMutex() throws Exception {
Semaphore sem = new Semaphore(1);
AtomicInteger cur = new AtomicInteger();
AtomicInteger peak = new AtomicInteger();
ExecutorService pool = Executors.newFixedThreadPool(4);
CountDownLatch done = new CountDownLatch(4);
for (int k = 0; k < 4; k++) {
pool.submit(() -> {
try {
sem.acquire();
try {
int c = cur.incrementAndGet();
peak.accumulateAndGet(c, Math::max);
sleep(30);
} finally {
cur.decrementAndGet();
sem.release();
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
} finally {
done.countDown();
}
return null;
});
}
done.await(10, TimeUnit.SECONDS);
pool.shutdown();
pool.awaitTermination(5, TimeUnit.SECONDS);
System.out.println("[2] peakConcurrent=" + peak.get() + " (same mutual exclusion as ReentrantLock)");
}
static void tryAcquireNoQueue() throws Exception {
Semaphore sem = new Semaphore(1);
sem.acquire();
System.out.println("[3] tryAcquireWhileHeld=" + sem.tryAcquire()
+ " availablePermits=" + sem.availablePermits());
System.out.println("[3] tryAcquireWithTimeout200ms=" + sem.tryAcquire(200, TimeUnit.MILLISECONDS)
+ " (waited 200ms in the queue, still nothing)");
sem.release();
System.out.println("[3] tryAcquireAfterRelease=" + sem.tryAcquire());
sem.release();
}
[1] peakConcurrent=3 availablePermits=3 queueLength=0
[2] peakConcurrent=1 (same mutual exclusion as ReentrantLock)
[3] tryAcquireWhileHeld=false availablePermits=0
[3] tryAcquireWithTimeout200ms=false (waited 200ms in the queue, still nothing)
[3] tryAcquireAfterRelease=true
peakConcurrent=3 说明限流生效,peakConcurrent=1 说明一许可信号量等价于互斥锁。第三组是面试常问的区别:tryAcquire() 无参版本不排队,拿不到立刻返回 false;tryAcquire(timeout) 会真的在队列里等满这段时间。两者配上 availablePermits() 就是一套非阻塞的限流判断——先试拿一个,拿不到就走降级逻辑,不占用线程。
Semaphore 的 state 语义和 CountDownLatch 正好相反:CountDownLatch 是"减到零放行",Semaphore 是"减到零阻塞"。acquire 减 state,减不动(remaining < 0)就返回负值让 AQS 去排队;release 加 state 并唤醒队列。所以 Semaphore 的许可可以还回去反复用,CountDownLatch 的计数归零就永久归零,这是它俩最本质的分工。
CyclicBarrier:能循环,但它根本不是 AQS
把需求换成"三个下游要齐步走,每轮都等最慢的那个,而且要复用",CountDownLatch 就不够用了——它计数归零即报废。这时该上 CyclicBarrier。
static void roundTrip() throws Exception {
AtomicInteger barrierActions = new AtomicInteger();
CyclicBarrier barrier = new CyclicBarrier(3, () -> {
barrierActions.incrementAndGet();
System.out.println(" round done");
});
AtomicInteger arrivals = new AtomicInteger();
ExecutorService pool = Executors.newFixedThreadPool(3);
for (int p = 0; p < 3; p++) {
final int id = p;
pool.submit(() -> {
for (int r = 0; r < 3; r++) {
sleep(20L * (id + 1)); // 故意让每方耗时不同
try {
barrier.await();
arrivals.incrementAndGet();
} catch (InterruptedException | BrokenBarrierException e) {
return null;
}
}
return null;
});
}
pool.shutdown();
pool.awaitTermination(5, TimeUnit.SECONDS);
System.out.println("[1] arrivals=" + arrivals.get() + " barrierActions=" + barrierActions.get()
+ " parties=" + barrier.getParties());
}
三方耗时分别是 20ms、40ms、60ms,故意不对称,看它是不是每轮都等齐。接着验证"等待人数"和"破损"这两个 CountDownLatch 没有的概念:
static void waitingCount() throws Exception {
CyclicBarrier barrier = new CyclicBarrier(3);
Thread t = new Thread(() -> {
try {
barrier.await();
} catch (Exception e) {
// 被 reset 打断
}
});
t.start();
while (barrier.getNumberWaiting() == 0) {
Thread.sleep(5);
}
System.out.println("[2] numberWaiting=" + barrier.getNumberWaiting() + " parties=" + barrier.getParties());
barrier.reset();
t.join();
}
static void brokenBarrier() throws Exception {
CyclicBarrier barrier = new CyclicBarrier(3);
AtomicInteger timeouts = new AtomicInteger();
AtomicInteger brokens = new AtomicInteger();
AtomicInteger oks = new AtomicInteger();
ExecutorService pool = Executors.newFixedThreadPool(3);
pool.submit(() -> { // 正常方:30ms 到场
sleep(30);
try { barrier.await(); oks.incrementAndGet(); }
catch (BrokenBarrierException e) { brokens.incrementAndGet(); }
catch (Exception e) { }
});
pool.submit(() -> { // 慢方:300ms 才到场
sleep(300);
try { barrier.await(); oks.incrementAndGet(); }
catch (BrokenBarrierException e) { brokens.incrementAndGet(); }
catch (Exception e) { }
});
pool.submit(() -> { // 超时方:只等 100ms
try { barrier.await(100, TimeUnit.MILLISECONDS); oks.incrementAndGet(); }
catch (TimeoutException e) { timeouts.incrementAndGet(); }
catch (BrokenBarrierException e) { brokens.incrementAndGet(); }
catch (Exception e) { }
});
pool.shutdown();
pool.awaitTermination(5, TimeUnit.SECONDS);
System.out.println("[3] timeout=" + timeouts.get() + " broken=" + brokens.get()
+ " ok=" + oks.get() + " isBroken=" + barrier.isBroken());
}
static void afterReset() throws Exception {
CyclicBarrier barrier = new CyclicBarrier(3);
barrier.reset();
AtomicInteger arrived = new AtomicInteger();
ExecutorService pool = Executors.newFixedThreadPool(3);
CountDownLatch done = new CountDownLatch(3);
for (int p = 0; p < 3; p++) {
pool.submit(() -> {
try {
barrier.await();
arrived.incrementAndGet();
} catch (Exception e) {
// 忽略
} finally {
done.countDown();
}
});
}
done.await(5, TimeUnit.SECONDS);
pool.shutdown();
pool.awaitTermination(5, TimeUnit.SECONDS);
System.out.println("[4] afterReset isBroken=" + barrier.isBroken() + " arrived=" + arrived.get());
}
round done
round done
round done
[1] arrivals=9 barrierActions=3 parties=3
[2] numberWaiting=1 parties=3
[3] timeout=1 broken=2 ok=0 isBroken=true
[4] afterReset isBroken=false arrived=3
四行输出四个结论。第一行组:三个 round done、barrierActions=3,说明三轮的屏障动作各触发一次,arrivals=9 说明九个 await() 全部返回——即使三方耗时 20/40/60ms 不等,每轮都被最慢的那方拖齐,屏障也没有被重复计数污染。
第二行 numberWaiting=1:getNumberWaiting() 是 CountDownLatch 完全没有的能力,它让外部能观测"有几个已经到场还在等"。
第三行是重点:timeout=1 broken=2。超时方抛 TimeoutException,同时它把整个屏障标记为破损,另外两方(包括还没到场的慢方)拿到的都是 BrokenBarrierException——一方出事,整批连坐。注意"慢方"的 300ms 其实还没到,等它到场时屏障已经破了,照样立刻拿异常。
第四行:reset() 之后 isBroken=false、arrived=3,屏障可以当新的用。reset() 同时会打破当前这一轮,把还在等的线程都用 BrokenBarrierException 放走。
现在说本文标题的那句话。CyclicBarrier 内部没有一个字段来自 AQS:
private static class Generation {
Generation() {} // prevent access constructor creation
boolean broken; // initially false
}
private final ReentrantLock lock = new ReentrantLock();
private final Condition trip = lock.newCondition();
private final int parties;
private final Runnable barrierCommand;
private Generation generation = new Generation();
private int count;
它用的是 ReentrantLock + Condition + 一个 Generation 代标记 + 一个普通的 int count。核心方法 dowait 一进门就 lock.lock(),然后 if (g.broken) throw new BrokenBarrierException();。所以它的状态载体是那个 int count(还没到场的方数),同步靠 AQS 的独占锁——ReentrantLock 本身是 AQS 独占模式,但 CyclicBarrier 不是 AQS 的子类,它只是个使用者。
Generation 这个类只有一个 boolean broken 字段,作用是把"第几轮"具象成一个对象:每轮结束 new 一个新的 Generation,旧的那个带着 broken 状态被丢弃。reset()、breakBarrier()、nextGeneration() 三个方法改的都是它。
Phaser:多阶段与动态注册
CyclicBarrier 有两个硬限制:参与方数量在构造时定死,且只有"全部到齐"一种推进条件。Phaser 把这两条都放开了。
static void phases() throws Exception {
Phaser phaser = new Phaser(3) {
@Override
protected boolean onAdvance(int phase, int registered) {
System.out.println(" onAdvance phase=" + phase + " registered=" + registered);
return phase == 2; // 第 3 阶段结束就终止
}
};
ExecutorService pool = Executors.newFixedThreadPool(3);
for (int p = 0; p < 3; p++) {
final int id = p;
pool.submit(() -> {
for (int round = 0; round < 3; round++) {
sleep(10L * (id + 1)); // 每阶段耗时不同
int returned = phaser.arriveAndAwaitAdvance();
if (id == 0) {
System.out.println(" party0 entered round " + round
+ " and got back " + returned);
}
}
return null;
});
}
pool.shutdown();
pool.awaitTermination(5, TimeUnit.SECONDS);
System.out.println("[1] terminated=" + phaser.isTerminated() + " phase=" + phaser.getPhase());
}
onAdvance(phase, registered) 是 Phaser 给出的钩子:每推进一个阶段调用一次,返回 true 表示"到此为止"。这里写 return phase == 2,也就是跑满三轮收工。它替代了 CyclicBarrier 那种"方数到零就自动推进"的固定规则——你可以按阶段号终止、按参与方数量终止、按业务条件终止。
另一半能力是动态注册和退出:
static void dynamic() throws Exception {
Phaser phaser = new Phaser(2); // 主线程 + 一个短工
Thread worker = new Thread(() -> {
sleep(30);
phaser.arriveAndDeregister(); // 到场并撤出,不再算一方
});
worker.start();
System.out.println("[2] before registeredParties=" + phaser.getRegisteredParties()
+ " arrived=" + phaser.getArrivedParties());
int phase = phaser.arriveAndAwaitAdvance();
System.out.println("[2] returnedFromAwait=" + phase
+ " registeredParties=" + phaser.getRegisteredParties()
+ " terminated=" + phaser.isTerminated());
phaser.arriveAndDeregister(); // 最后一方也撤出
System.out.println("[2] afterAllLeave terminated=" + phaser.isTerminated());
worker.join();
}
onAdvance phase=0 registered=3
party0 entered round 0 and got back 1
onAdvance phase=1 registered=3
party0 entered round 1 and got back 2
onAdvance phase=2 registered=3
party0 entered round 2 and got back -2147483645
[1] terminated=true phase=-2147483645
[2] before registeredParties=2 arrived=0
[2] returnedFromAwait=1 registeredParties=1 terminated=false
[2] afterAllLeave terminated=true
arriveAndAwaitAdvance() 的返回值是最容易搞错的地方:它返回推进之后的新阶段号,不是旧的那个。party0 三轮分别拿到 1、2、-2147483645,前两个就是新阶段号的真身,第三个是终止信号。
-2147483645 这个数字不是随便来的。Phaser 的内部 state 是一个 volatile long,高位存阶段号、中间存参与方、低位存未到场数,常量定义是:
MAX_PARTIES = 0xffff
MAX_PHASE = Integer.MAX_VALUE
PARTIES_SHIFT = 16
PHASE_SHIFT = 32
UNARRIVED_MASK = 0xffff
PARTIES_MASK = 0xffff0000L
COUNTS_MASK = 0xffffffffL
TERMINATION_BIT = 1L << 63
ONE_ARRIVAL = 1
ONE_PARTY = 1 << PARTIES_SHIFT
ONE_DEREGISTER = ONE_ARRIVAL | ONE_PARTY
EMPTY = 1
推进阶段的那段代码是 if (onAdvance(phase, nextUnarrived)) n |= TERMINATION_BIT;,紧接着 int nextPhase = (phase + 1) & MAX_PHASE;。TERMINATION_BIT = 1L << 63 落在整个 long 的最高位,换算到"阶段号"这个 32 位字段上正好是符号位。所以第三轮 onAdvance 返回 true 之后,新阶段号是 3,再或上终止位,getPhase() 把高 32 位当有符号 int 读出来就是 0x80000003,也就是 -2147483645。把 -2147483645 换算回去:加满 2³² 得 2147483651,十六进制 0x80000003,最高位是终止标记,低两位是十进制 3——和"跑完三轮"对得上。
[2] 那三行是动态注册的证据:一开始 registeredParties=2,短工 arriveAndDeregister() 之后变成 1,主线程不用改任何配置就继续推进;最后一方也退出,terminated=true。Phaser 还允许后来者 register() 加入,这是 CyclicBarrier 做不到的。
和 CyclicBarrier 一样,Phaser 也不基于 AQS。它自己的 state 是 volatile long,靠 VarHandle(源码里的 STATE)做 CAS;排队节点是 QNode implements ForkJoinPool.ManagedBlocker——注意这个接口,它意味着在 ForkJoinPool 里等待时可以让出工作线程去干别的活,这是 AQS 的 CLH 队列给不了的。
Exchanger:一次只换一对
家族图右栏还剩最后一个。Exchanger 解决的不是"等齐"也不是"限流",而是两个线程在一个汇合点上互换手上的数据:exchange(V) 会一直阻塞,直到另一个线程也调用它,然后双方各自拿到对方传进来的那个对象返回。一次只配一对,多个线程同时来就按到达顺序两两配。
import java.util.concurrent.Exchanger;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
/**
* Exchanger 的三件事:
* 1. 两个线程各带一份数据到汇合点,互换之后各拿对方的
* 2. 对方一直不来,带超时的 exchange 抛 TimeoutException
* 3. 同一个 Exchanger 可以反复用,同一对线程能换多轮
* 用法: java ExchangerDemo
*/
public class ExchangerDemo {
public static void main(String[] args) throws Exception {
Exchanger<String> ex = new Exchanger<>();
StringBuilder workerLog = new StringBuilder();
Thread worker = new Thread(() -> {
try {
workerLog.append("worker got=").append(ex.exchange("worker-parcel"));
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
});
worker.start();
String mainGot = ex.exchange("main-parcel");
worker.join();
System.out.println("[1] mainGot=" + mainGot + " " + workerLog);
Exchanger<String> solo = new Exchanger<>();
try {
solo.exchange("x", 200, TimeUnit.MILLISECONDS);
System.out.println("[2] returned normally (unexpected)");
} catch (TimeoutException e) {
System.out.println("[2] TimeoutException after 200ms (no partner arrived)");
}
Exchanger<String> reuse = new Exchanger<>();
StringBuilder rounds = new StringBuilder();
Thread w2 = new Thread(() -> {
try {
String r1 = reuse.exchange("w-round1");
String r2 = reuse.exchange("w-round2");
rounds.append("worker got ").append(r1).append(" then ").append(r2);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
});
w2.start();
String m1 = reuse.exchange("m-round1");
String m2 = reuse.exchange("m-round2");
w2.join();
System.out.println("[3] mainGot=" + m1 + " then " + m2 + " " + rounds);
}
}
[1] mainGot=worker-parcel worker got=main-parcel
[2] TimeoutException after 200ms (no partner arrived)
[3] mainGot=w-round1 then w-round2 worker got m-round1 then m-round2
三行输出三件事。[1] 两边各自把自己那份传进去、把对方那份拿回来;谁先到都行,先到的那个阻塞等着,配对成功才一起返回。[2] 说明它也有带超时的版本,exchange(V, timeout, unit) 在没人来的时候抛 TimeoutException 而不是永久挂着——接不上就降级,是它唯一安全的用法。[3] 验证它可以反复使用:同一个 Exchanger 上同一对线程连换两轮,第一轮的配对不会污染第二轮。
它和 CyclicBarrier 的分工也在这里:屏障只要求"都到齐",不留数据;Exchanger 到齐之后还要交换数据,所以它天然只能两两配对。典型场景只有一个形状——两条流水线在某个点对调缓冲,比如一条线程读、一条线程写,缓冲区轮换使用,省掉一次拷贝。
Exchanger 同样不基于 AQS:它自己维护一个 Slot[] 槽位数组和一个 volatile int bound(靠 VarHandle 做 CAS),核心逻辑收在私有的 xchg(V, long) 里。
ReentrantReadWriteLock:state 被切成两半,切法换过
读多写少的场景适合读写锁。ReentrantReadWriteLock 要在一个 state 里同时装下读计数、写重入次数、写锁归属,做法是把 state 切成两半。这一段先看它继承谁:
static void parentClass() throws Exception {
Class<?> sync = Class.forName("java.util.concurrent.locks.ReentrantReadWriteLock$Sync");
System.out.println("[1] Sync extends " + sync.getSuperclass().getName());
}
用 Class.forName 直接问 JVM,比翻文档可靠——这个类是包级私有的,只能反射拿到。
上面这张图是本文最需要注意版本差异的地方。JDK 8 到 25 的资料会告诉你"高 16 位记读锁、低 16 位记写锁",因为那时 Sync 继承的是 AbstractQueuedSynchronizer,state 是 32 位 int,只能对半切成 16+16。JDK 26 里 Sync 已经改成继承 AbstractQueuedLongSynchronizer,state 变成 64 位 long,切法变成 32+32:
SHARED_SHIFT = 16 SHARED_UNIT = 1 << 16 = 65536 MAX_COUNT = 65535 EXCLUSIVE_MASK = 0xFFFF
上面这行是 JDK 8~25 的常量,下面是 JDK 26 的:
SHARED_SHIFT = 32 SHARED_UNIT = 1L << 32 MAX_COUNT = 2147483647 EXCLUSIVE_MASK = (1L << 32) - 1
对应的源码声明可以直接翻到,在 ReentrantReadWriteLock.java 第 258 行起:
abstract static class Sync extends AbstractQueuedLongSynchronizer {
static final int SHARED_SHIFT = 32;
static final long SHARED_UNIT = (1L << SHARED_SHIFT);
static final long MAX_COUNT = Integer.MAX_VALUE;
static final long EXCLUSIVE_MASK = (1L << SHARED_SHIFT) - 1;
static int sharedCount(long c) { return (int)(c >>> SHARED_SHIFT); }
static int exclusiveCount(long c) { return (int)(c & EXCLUSIVE_MASK); }
}
方法名一个字没改,sharedCount 还是那个 sharedCount,只有右移的位数从 16 变成了 32。改动的动机很实在:16 位一半最多放 65535,读锁持有数或写重入次数一超就抛 Error("Maximum lock count exceeded"),换成 long 之后上限提到 21 亿(这个改动对应 OpenJDK 的 JDK-8352971)。
如果你手上的 JDK 不是 26,可以用一条命令确认自己的版本属于哪一档:
javap -constants "java.util.concurrent.locks.ReentrantReadWriteLock$Sync"
双引号不能省,$Sync 会被 PowerShell 当成变量展开。
切法搞清楚了,再看行为。下面这段依次验证读并行、写互斥、持读锁抢写锁、写锁降级、读锁重入:
static void readParallelWriteExclusive() throws Exception {
ReentrantReadWriteLock rw = new ReentrantReadWriteLock();
CountDownLatch acquired = new CountDownLatch(4);
CountDownLatch release = new CountDownLatch(1);
ExecutorService pool = Executors.newFixedThreadPool(4);
for (int k = 0; k < 4; k++) {
pool.submit(() -> {
rw.readLock().lock();
try {
acquired.countDown();
release.await();
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
} finally {
rw.readLock().unlock();
}
return null;
});
}
acquired.await(5, TimeUnit.SECONDS);
System.out.println("[2] readLockCount=" + rw.getReadLockCount()
+ " writeLocked=" + rw.isWriteLocked() + " queuedWriters=" + rw.getQueueLength());
release.countDown();
pool.shutdown();
pool.awaitTermination(5, TimeUnit.SECONDS);
// 两个写线程:第一个拿到锁不放,第二个只能排队
CountDownLatch first = new CountDownLatch(1);
CountDownLatch releaseFirst = new CountDownLatch(1);
ExecutorService pool2 = Executors.newFixedThreadPool(2);
pool2.submit(() -> {
rw.writeLock().lock();
try {
first.countDown();
releaseFirst.await();
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
} finally {
rw.writeLock().unlock();
}
return null;
});
first.await(5, TimeUnit.SECONDS);
pool2.submit(() -> {
boolean tryBoth = rw.writeLock().tryLock();
System.out.println(" secondWriterTryLock=" + tryBoth);
rw.writeLock().lock();
try {
System.out.println(" secondWriterGotInAfterRelease=true");
} finally {
rw.writeLock().unlock();
}
return null;
});
Thread.sleep(200);
System.out.println("[3] whileFirstWriterHolds writeLocked=" + rw.isWriteLocked()
+ " queuedWriters=" + rw.getQueueLength());
releaseFirst.countDown();
pool2.shutdown();
pool2.awaitTermination(5, TimeUnit.SECONDS);
}
static void upgradeSelfDeadlock() throws Exception {
ReentrantReadWriteLock rw = new ReentrantReadWriteLock();
Thread t = new Thread(() -> {
rw.readLock().lock();
try {
System.out.println(" upgradeTryLock=" + rw.writeLock().tryLock());
System.out.println(" upgradeTryLock200ms="
+ rw.writeLock().tryLock(200, TimeUnit.MILLISECONDS));
rw.writeLock().lock(); // 永远等不到:自己就是那个读者
System.out.println(" never printed");
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
} finally {
rw.readLock().unlock();
}
});
t.setDaemon(true);
t.start();
Thread.sleep(500);
System.out.println("[4] upgraderState=" + t.getState() + " stillAlive=" + t.isAlive()
+ " (its own read lock is never released, so the write lock never arrives)");
}
static void downgrade() {
ReentrantReadWriteLock rw = new ReentrantReadWriteLock();
rw.writeLock().lock();
System.out.println("[5] holdingWrite writeLocked=" + rw.isWriteLocked()
+ " writeHoldCount=" + rw.getWriteHoldCount());
rw.readLock().lock(); // 持写锁时拿读锁:允许
rw.writeLock().unlock(); // 先放写,再当读者
System.out.println("[5] afterDowngrade writeLocked=" + rw.isWriteLocked()
+ " readLockCount=" + rw.getReadLockCount()
+ " readHoldCount=" + rw.getReadHoldCount());
rw.readLock().unlock();
System.out.println("[5] released readLockCount=" + rw.getReadLockCount());
}
static void readReentrant() {
ReentrantReadWriteLock rw = new ReentrantReadWriteLock();
rw.readLock().lock();
rw.readLock().lock();
System.out.println("[6] readHoldCount=" + rw.getReadHoldCount()
+ " readLockCount=" + rw.getReadLockCount());
rw.readLock().unlock();
rw.readLock().unlock();
}
[1] Sync extends java.util.concurrent.locks.AbstractQueuedLongSynchronizer
[2] readLockCount=4 writeLocked=false queuedWriters=0
secondWriterTryLock=false
[3] whileFirstWriterHolds writeLocked=true queuedWriters=1
secondWriterGotInAfterRelease=true
upgradeTryLock=false
upgradeTryLock200ms=false
[4] upgraderState=WAITING stillAlive=true (its own read lock is never released, so the write lock never arrives)
[5] holdingWrite writeLocked=true writeHoldCount=1
[5] afterDowngrade writeLocked=false readLockCount=1 readHoldCount=1
[5] released readLockCount=0
[6] readHoldCount=2 readLockCount=2
第一行回答了"它站在谁的肩膀上":JDK 26 里 Sync 的父类就是 AbstractQueuedLongSynchronizer,不用再靠记忆。
[2] 里四个读者同时持锁,readLockCount=4 且 writeLocked=false——读锁互相不挡。[3] 里写锁一旦被持有,secondWriterTryLock=false,第二个写线程进队列(queuedWriters=1),释放后才拿到。
[4] 是这一段最该记住的坑:upgradeTryLock=false、upgradeTryLock200ms=false,然后 rw.writeLock().lock() 把线程永久挂住,upgraderState=WAITING stillAlive=true。注意这里只有一个线程——它自己拿着读锁,又去申请写锁。写锁的前提是"读者数为零",而它自己就是那个读者,于是自己等自己。很多资料写的是"两个线程同时持读锁再抢写锁会死锁",这话没错但抓错了重点:两个线程的情况至少还能靠双方都主动释放读锁来解开,单线程自我升级是逻辑上不可解的,只能靠 tryLock 提前发现。所以升级必须写成先放读锁再拿写锁,中间的重检查不能省。
不过 [4] 有个细节值得说清楚:tryLock() 失败不是因为"读锁卡住了写锁",而是因为写锁的前提条件没满足。看 JDK 26 的 tryAcquireShared(读锁的获取逻辑,写锁走它的镜像 tryAcquire)第一段:
long c = getState();
if (exclusiveCount(c) != 0 &&
getExclusiveOwnerThread() != current)
return -1L;
exclusiveCount(c) != 0 表示有写锁被持有;只有当持写锁的不是当前线程时才返回 -1 去排队。反过来,如果写锁已经在自己手里,这个判断直接放行,读锁就能拿到——这正是 [5] 里降级能成功的原因。源码在 tryAcquire 里对这个分支有一句注释:else we hold the exclusive lock; blocking here would cause deadlock(否则我们正持有独占锁,在这里阻塞会导致死锁)。所以降级的真正原因是当前线程就是写锁持有者,判断条件放它过去,而不是"读锁的 readerShouldBlock 对它放行"——后者管的是公平性,是另一回事。
[5] 的输出 holdingWrite writeLocked=true → afterDowngrade writeLocked=false readLockCount=1 readHoldCount=1:写锁放掉、读锁留着,这正是降级的用途——改完数据立刻以读者身份继续用,中间不留一个"谁都不持有"的空档。[6] 的 readHoldCount=2 readLockCount=2 说明读锁可重入,而且内外两个计数分别记录(总读者数 vs 本线程持有数)。
StampedLock:stamp 就是 state 的快照
StampedLock 是这三把锁里唯一"不用锁"的——它把乐观读做成了 API。先把它的 state 布局和乐观读时序看清楚:
内部 state 是一个 long,切成三段:低 7 位是读锁持有数,第 8 位是写锁位,高 56 位是版本号。常量定义在 StampedLock.java 第 314 行起:
LG_READERS = 7 RUNIT = 1L WBIT = 1L << 7 = 128 RBITS = WBIT - 1 = 127 ORIGIN = WBIT << 1 = 256
WBIT = 1L << 7 就是 0x80,正好落在第 8 位;RBITS = 127 是低 7 位的全 1 掩码;ORIGIN = 256 是新建锁时的初始 state(版本号从 1 开始,低 8 位全零表示无锁)。下面这段把它跑出来看:
static void stampLayout() {
StampedLock lock = new StampedLock();
long optimistic = lock.tryOptimisticRead(); // 没锁的时候返回 state 的版本部分
long write = lock.writeLock();
lock.unlockWrite(write); // 释放写锁会把版本号 +1
long read = lock.readLock();
long read2 = lock.readLock();
int holding = lock.getReadLockCount();
lock.unlockRead(read);
lock.unlockRead(read2);
long optimisticAgain = lock.tryOptimisticRead();
System.out.printf("[1] optimistic=0x%X write=0x%X read=0x%X%n", optimistic, write, read);
System.out.printf("[1] low8 optimistic=0x%02X write=0x%02X read=0x%02X%n",
optimistic & 0xFF, write & 0xFF, read & 0xFF);
System.out.println("[1] readLockCountWhileTwoHeld=" + holding);
System.out.printf("[1] optimisticAfterOneWriteCycle=0x%X%n", optimisticAgain);
}
[1] optimistic=0x100 write=0x180 read=0x201
[1] low8 optimistic=0x00 write=0x80 read=0x01
[1] readLockCountWhileTwoHeld=2
[1] optimisticAfterOneWriteCycle=0x200
把四个数摊开看:新建后 0x100(版本 1、低 8 位全零,无锁),取写锁变成 0x180(第 8 位置 1),释放写锁后版本号加一变成 0x200,再取一把读锁变成 0x201(低 7 位记了 1 个读者)。low8 那一行是同一件事的验证:不带锁时低 8 位是 0x00,写锁是 0x80,读锁是 0x01——低 8 位就是模式位。所以 tryOptimisticRead() 返回的压根不是锁,是"当时 state 去掉低 7 位"的那份快照。
乐观读的完整用法要配合 validate。这段分别制造"写已经发生"和"写正在进行"两种情况:
static void staleRead() throws Exception {
StampedLock lock = new StampedLock();
int[] value = {1};
long stamp = lock.tryOptimisticRead();
int seen = value[0];
Thread writer = new Thread(() -> {
long ws = lock.writeLock();
try {
value[0] = 2;
} finally {
lock.unlockWrite(ws);
}
});
writer.start();
writer.join();
boolean ok = lock.validate(stamp);
System.out.println("[2] optimisticRead=" + seen + " validateAfterWrite=" + ok);
long rs = lock.readLock();
try {
seen = value[0];
} finally {
lock.unlockRead(rs);
}
long fresh = lock.tryOptimisticRead();
System.out.println("[2] retryRead=" + seen + " validateFresh=" + lock.validate(fresh));
}
static void validateWhileWriting() throws Exception {
StampedLock lock = new StampedLock();
long stamp = lock.tryOptimisticRead();
CountDownLatch holding = new CountDownLatch(1);
CountDownLatch release = new CountDownLatch(1);
Thread writer = new Thread(() -> {
long ws = lock.writeLock();
holding.countDown();
try {
release.await();
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
} finally {
lock.unlockWrite(ws);
}
});
writer.start();
holding.await();
boolean whileWriting = lock.validate(stamp);
release.countDown();
writer.join();
boolean afterWrite = lock.validate(stamp);
System.out.println("[3] whileWriting=" + whileWriting + " afterWriteSameStamp=" + afterWrite);
}
[2] optimisticRead=1 validateAfterWrite=false
[2] retryRead=2 validateFresh=true
[3] whileWriting=false afterWriteSameStamp=false
[2] 是"写已经发生":乐观读到旧值 1,写线程改完并释放后 validateAfterWrite=false,于是老老实实加读锁重读,拿到 2,重取的新快照验证通过(validateFresh=true)。[3] 是"写正在进行":同一个 stamp 在写锁被持有期间验一次、写锁释放后再验一次,false 和 false——两种情况下 validate 都给 false。
为什么两种都失败?看 validate 的实现,JDK 26 里就两行:
public boolean validate(long stamp) {
U.loadFence();
return (stamp & SBITS) == (state & SBITS);
}
SBITS = ~RBITS,把低 7 位屏蔽掉,只比高位。写已经发生——版本号从 V 变成 V+1,高位不等;写正在进行——第 8 位(WBIT)被置 1,也在高位里,同样不等。所以乐观读挡不住的东西比想象中多:它不是"读的时候没别人写",而是"读完再验一次,中间有人写过就重来"。验不过就必须重读,拿旧值硬用等于自己制造脏读。
剩下两个限制也顺便验掉:
static void notReentrant() {
StampedLock lock = new StampedLock();
long ws = lock.writeLock();
long again = lock.tryWriteLock();
long readWhileWriting = lock.tryReadLock();
System.out.println("[4] holdingWrite tryWriteLockAgain=" + again
+ " tryReadLock=" + readWhileWriting + " (both refused)");
lock.unlockWrite(ws);
long ok = lock.tryWriteLock();
System.out.println("[4] afterUnlock tryWriteLock=" + (ok != 0L));
if (ok != 0L) {
lock.unlockWrite(ok);
}
}
static void convertModes() {
StampedLock lock = new StampedLock();
long rs = lock.readLock();
long ws = lock.tryConvertToWriteLock(rs);
System.out.println("[5] readToWrite=" + (ws != 0L) + " writeLocked=" + lock.isWriteLocked());
long back = lock.tryConvertToReadLock(ws);
System.out.println("[5] writeToRead=" + (back != 0L) + " writeLocked=" + lock.isWriteLocked()
+ " readLockCount=" + lock.getReadLockCount());
lock.unlockRead(back);
}
static void interruptibleVariants() throws Exception {
System.out.println("[6] " + StampedLock.class.getMethod("readLockInterruptibly"));
System.out.println("[6] " + StampedLock.class.getMethod("writeLockInterruptibly"));
}
[4] holdingWrite tryWriteLockAgain=0 tryReadLock=0 (both refused)
[4] afterUnlock tryWriteLock=true
[5] readToWrite=true writeLocked=true
[5] writeToRead=true writeLocked=false readLockCount=1
[4] 验证不可重入:已经拿着写锁,再 tryWriteLock 返回 0,连 tryReadLock 也返回 0——注意 StampedLock 的约定是失败返回 0,不是抛异常也不是返回 false,判断时要写 != 0L 而不是直接当布尔用。释放之后立刻能拿到。
[5] 是 StampedLock 相对读写锁的一个真实优势:只有一个读者时,tryConvertToWriteLock 直接把读锁升级成写锁且返回非 0(readToWrite=true);tryConvertToReadLock 反向也成立。前面 ReentrantReadWriteLock 那节刚证明它的读锁升级是自锁的,这里是同一个问题的另一种解法——但 tryConvertToWriteLock 只在"恰好一个读者"时成功,多一个读者就返回 0,所以它治的是单线程场景,不是通用的升级方案。
[6] 是打印 Method 对象本身,用来一次性澄清一个常见说法。"StampedLock 中断也不响应"这句话不准确:无参的 readLock()、writeLock() 确实不响应中断,但 readLockInterruptibly() 和 writeLockInterruptibly() 是专门提供的,需要中断语义时用它们。
线程六态:BLOCKED 只属于 synchronized
前面已经见过 WAITING 了——[4] 里永久的 upgraderState=WAITING 就是。六个状态的名字来自 Thread.State 枚举,迁移关系如下:
真正需要记的是三个"在等"的状态之间的分工:BLOCKED 只等 synchronized 的监视器,WAITING 是没有超时地等唤醒,TIMED_WAITING 是带超时地等。这三个名字在 jstack 输出里原样出现,认状态比认堆栈更快。跑一段代码逐个确认:
static volatile boolean spin = true;
static void waitUntil(Thread t, Thread.State want) throws Exception {
while (t.getState() != want) {
Thread.sleep(2);
}
}
/** NEW -> RUNNABLE -> TERMINATED。 */
static void basic() throws Exception {
CountDownLatch started = new CountDownLatch(1);
Thread t = new Thread(() -> {
started.countDown();
while (spin) {
Thread.onSpinWait();
}
});
System.out.println("[1] beforeStart=" + t.getState());
t.start();
started.await();
Thread.sleep(50);
System.out.println("[1] whileSpinning=" + t.getState());
spin = false;
t.join();
System.out.println("[1] afterJoin=" + t.getState());
}
/** BLOCKED 只属于 synchronized 监视器。 */
static void blocked() throws Exception {
Object monitor = new Object();
CountDownLatch held = new CountDownLatch(1);
CountDownLatch release = new CountDownLatch(1);
Thread holder = new Thread(() -> {
synchronized (monitor) {
held.countDown();
try {
release.await();
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
});
holder.start();
held.await();
Thread blocker = new Thread(() -> {
synchronized (monitor) {
// 拿到就退出
}
});
blocker.start();
waitUntil(blocker, Thread.State.BLOCKED);
System.out.println("[2] waitingForMonitor=" + blocker.getState());
release.countDown();
blocker.join();
holder.join();
}
/** WAITING 三兄弟:park、wait、join,全都没有超时。 */
static void waiting() throws Exception {
Object lock = new Object();
CountDownLatch parked = new CountDownLatch(1);
CountDownLatch waited = new CountDownLatch(1);
CountDownLatch joinStarted = new CountDownLatch(1);
CountDownLatch gate = new CountDownLatch(1);
Thread p = new Thread(() -> {
parked.countDown();
LockSupport.park();
});
Thread w = new Thread(() -> {
waited.countDown();
synchronized (lock) {
try {
lock.wait();
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
});
Thread j = new Thread(() -> {
joinStarted.countDown();
try {
gate.await();
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
});
p.start(); w.start(); j.start();
parked.await(); waited.await(); joinStarted.await();
waitUntil(p, Thread.State.WAITING);
waitUntil(w, Thread.State.WAITING);
waitUntil(j, Thread.State.WAITING);
System.out.println("[3] park=" + p.getState() + " wait=" + w.getState()
+ " latchAwait=" + j.getState());
LockSupport.unpark(p);
synchronized (lock) {
lock.notifyAll();
}
gate.countDown();
p.join(); w.join(); j.join();
System.out.println("[3] allDone");
}
/** TIMED_WAITING:带超时的那一版。 */
static void timedWaiting() throws Exception {
CountDownLatch ready = new CountDownLatch(4);
Object lock = new Object();
Thread sleepThread = new Thread(() -> {
ready.countDown();
sleep(10_000);
});
Thread parkNanosThread = new Thread(() -> {
ready.countDown();
LockSupport.parkNanos(10_000_000_000L);
});
Thread waitThread = new Thread(() -> {
ready.countDown();
synchronized (lock) {
try {
lock.wait(10_000);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
});
Thread anchor = new Thread(() -> sleep(10_000));
anchor.setDaemon(true);
anchor.start();
Thread joinThread = new Thread(() -> {
ready.countDown();
try {
anchor.join(10_000);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
});
sleepThread.start(); parkNanosThread.start(); waitThread.start(); joinThread.start();
ready.await();
waitUntil(sleepThread, Thread.State.TIMED_WAITING);
waitUntil(parkNanosThread, Thread.State.TIMED_WAITING);
waitUntil(waitThread, Thread.State.TIMED_WAITING);
waitUntil(joinThread, Thread.State.TIMED_WAITING);
System.out.println("[4] sleep=" + sleepThread.getState()
+ " parkNanos=" + parkNanosThread.getState()
+ " waitN=" + waitThread.getState()
+ " joinN=" + joinThread.getState());
sleepThread.interrupt();
parkNanosThread.interrupt();
Thread.sleep(50);
synchronized (lock) {
lock.notifyAll();
}
joinThread.interrupt();
sleepThread.join(); parkNanosThread.join(); waitThread.join(); joinThread.join();
}
[1] beforeStart=NEW
[1] whileSpinning=RUNNABLE
[1] afterJoin=TERMINATED
[2] waitingForMonitor=BLOCKED
[3] park=WAITING wait=WAITING latchAwait=WAITING
[3] allDone
[4] sleep=TIMED_WAITING parkNanos=TIMED_WAITING waitN=TIMED_WAITING joinN=TIMED_WAITING
[3] 那一行信息量最大:LockSupport.park()、Object.wait()、CountDownLatch.await() 三种完全不同的实现,线程状态都是 WAITING。这就是"AQS 用 LockSupport.park 阻塞"的可见证据——synchronized 抢不到监视器才是 BLOCKED,AQS 队列上的等待线程在 jstack 里看到的是 WAITING。
waitUntil 这个方法本身也值得留意:它用"轮询到目标状态"代替 Thread.sleep(100) 这类固定等待。并发演示要稳定,就不能赌"100ms 应该够了"——机器一忙就飘。所有演示输出能两次跑出同一份结果,靠的都是"用 latch 串时序 + 轮询状态",而不是睡眠时长。
这一段还有个坑值得提前说:spin 必须声明成 static volatile boolean。写成局部数组 boolean[] spin = {true} 让 lambda 捕获,数组元素本身不是 volatile,JIT 完全可以把读操作提到循环外面,工作线程就永远退不出来了。这个"看着能跑、偶尔挂死"的 bug 在 visibility 类演示里非常常见。
中断、park 与许可:三件容易答反的事
上一节那段代码里还有两行输出没说,它们对应中断语义最容易答错的地方:
parkReturned interruptedFlag=true
sleepThrew=InterruptedException interruptedFlagAfterCatch=false
同一个线程被 interrupt(),LockSupport.park() 的做法是静默返回,标志位保持为 true(parkReturned interruptedFlag=true);而 Thread.sleep() 的做法是抛 InterruptedException,同时把标志位清掉(sleepThrew=InterruptedException interruptedFlagAfterCatch=false)。这两行是背不出来的,必须跑出来看。
生产含义很直接:park 醒来之后你必须自己判断"是正常唤醒还是被中断了",用 Thread.currentThread().isInterrupted() 检查;sleep 那些方法则要求你在 catch 里决定要不要 Thread.currentThread().interrupt() 把标志位补回去——不补,上游就永远不知道有人想停这个线程。
后面还有两行:
[5] firstParkReturnedUnder5ms=true parkNanosWaitedAtLeast15ms=true
[6] blockerVisibleToOthers=true
[5] 验证 park 的许可机制。代码是连续两次 unpark 再 park:第一个 park 立刻返回(firstParkReturnedUnder5ms=true),因为许可可以提前发放;第二个许可以及后续的 parkNanos(20ms) 就没人管了,只能等满时间(parkNanosWaitedAtLeast15ms=true)。所以 LockSupport 的许可是"最多攒一个"的计数器,多给的会丢。这是它比 wait/notify 好的地方——notify 在没人等的时候调用,信号直接丢掉;unpark 会先存着,下次 park 立刻消费。
[6] 验证 LockSupport.getBlocker(t) 能拿到阻塞原因对象。调试线上"线程卡在哪"时,jstack 里 park 的线程会显示这个对象——把锁对象、超时时间塞进去,排查成本立刻降下来。
最后回到 Semaphore 的一个特例,它把上面这些机制串在了一起:
static void uninterruptible() throws Exception {
Semaphore sem = new Semaphore(1);
CountDownLatch queued = new CountDownLatch(1);
sem.acquire(); // 主线程占住唯一许可
Thread t = new Thread(() -> {
queued.countDown();
sem.acquireUninterruptibly(); // 排队等
System.out.println(" acquireUninterruptibly returned, interruptedFlag="
+ Thread.currentThread().isInterrupted());
});
t.start();
queued.await();
Thread.sleep(100);
t.interrupt();
Thread.sleep(100);
System.out.println("[4] afterInterrupt state=" + t.getState()
+ " interruptedFlag=" + t.isInterrupted() + " availablePermits=" + sem.availablePermits());
sem.release(); // 放许可,它才继续
t.join();
}
[4] afterInterrupt state=WAITING interruptedFlag=false availablePermits=0
acquireUninterruptibly returned, interruptedFlag=true
这两行是全文最反直觉的一段。中断之后,排队线程的状态是 WAITING(没有被"叫醒",因为 acquireUninterruptibly 不响应中断),而且 interruptedFlag=false——标志位此刻是 false。等到许可释放、它真正拿到许可之后,打印出来的 interruptedFlag=true。标志位在被中断时被清了一下,然后又被补回去了。
原因是 AQS 内部的排队逻辑:等待过程中它调用的是会清标志位的中断检查,把"被中断过"这件事记在一个局部变量里;等真正拿到资源、成为队列头部之后,再把这个标志补回去。JDK 26 里这段逻辑在统一的 acquire(node, arg, shared, interruptible, timed, time) 方法中,成功分支写的是 if (interrupted) current.interrupt();。
这里又是一个版本差异,引用方法名时要小心:JDK 8 的 AQS 里对应的是 doAcquireShared(int) 加 selfInterrupt();从 JDK 14 起已经换成上面这个统一的 acquire(...),JDK 26 沿用。所以"查 AQS 源码确认 doAcquireShared 怎么写的"这种做法,在新 JDK 上会直接搜不到方法。
到底怎么选
七个类、五种 state 设计看下来,选型可以收敛成一张表:
| 工具 | state 是什么 | 同步基础 | 一次性/可复用 | 适合的场景 |
|---|---|---|---|---|
CountDownLatch | AQS int state = 剩余计数 | AQS 共享模式 | 一次性,归零即废 | 一个线程等 N 个任务齐活 |
Semaphore | AQS int state = 可用许可 | AQS 共享模式 | 可复用 | 限流、连接池、一许可当互斥锁 |
CyclicBarrier | 自己的 int count + Generation | ReentrantLock + Condition | 可复用,但会破损 | N 个线程多轮齐步走 |
Phaser | 自己的 volatile long state | VarHandle CAS + QNode | 可复用,可动态增删 | 多阶段 + 参与方数量会变 |
Exchanger | 自己的 Slot[] + volatile int bound | VarHandle CAS | 可复用 | 两条流水线对调缓冲,一次配一对 |
ReentrantReadWriteLock | AQLS long state 32+32 对半切 | AQLS 共享 + 独占 | 可复用 | 读多写少且需要公平/条件变量 |
StampedLock | 自己的 long:低 7 位读、第 8 位写、高位版本 | CAS + 乐观读 | 可复用 | 读极多写极少,且能接受不可重入 |
按这张表往下推,几个判断点是有先后顺序的:
先问"等的是计数还是批次"。计数归零就结束、之后不再用,选 CountDownLatch;要反复等齐、每轮重新计数,往 CyclicBarrier 或 Phaser 看。
再问"参与方数量会变吗、阶段条件复杂吗"。会变、或者终止条件不只是"到齐",选 Phaser;固定 N 方、只要到齐,CyclicBarrier 更简单。Phaser 的 onAdvance 钩子是 CyclicBarrier 给不了的。
然后问"限流还是同步"。控制并发数量用 Semaphore,它是唯一一个"持有期间可以不释放、按需还回去"的工具。
再问"要不要交换数据"。只是等齐就用屏障,等齐之后还要把数据换过去,才轮到 Exchanger——它是这张表里唯一一个会"返回对方给你的东西"的工具,也因此只能两两配对。
最后问"读写比例"。读写锁和 StampedLock 都要求读多写少,但代价不同:ReentrantReadWriteLock 支持重入、支持 Condition、能公平,升级会自锁;StampedLock 用乐观读换吞吐,代价是不可重入、没有 Condition、validate 失败要自己重试,而且它的锁不支持 ReentrantLock 那样的可中断版本(要用 readLockInterruptibly / writeLockInterruptibly)。读操作特别密集、写极少、并且你能接受回退逻辑,才轮到 StampedLock。
回到面试速答,把这段最容易答错的四条重新写一遍:
- 这一族里基于 AQS 的只有三个:
CountDownLatch和Semaphore用 AQS 共享模式,ReentrantLock用 AQS 独占模式。CyclicBarrier是ReentrantLock+Condition+Generation,Phaser和StampedLock各自维护volatile long state,Exchanger是Slot[]加bound——它们都不是 AQS 子类。 ReentrantReadWriteLock的 state 切分要带版本号:JDK 8~25 是int的 16+16,JDK 26 起是long的 32+32。读锁不可升级,单线程持读锁抢写锁就是自锁;写锁可降级,原因是tryAcquire里"写锁被持有且持有者不是当前线程才返回 -1",源码注释写明在此阻塞会导致死锁。StampedLock的 state 是低 7 位读计数、第 8 位写锁位、高 56 位版本号;validate比的是(stamp & SBITS) == (state & SBITS),只比高位,所以"写已经发生"和"写正在进行"都会失败。不可重入,失败返回 0 而不是抛异常。- 线程六态里
BLOCKED只对应synchronized监视器;LockSupport.park()是WAITING,park被中断是静默返回且标志位保留,sleep/wait/join是抛异常且标志位被清。park的许可最多攒一个,可以提前发放。
核心收获就一句:这些工具类的差异不是 API 风格差异,而是 state 语义的差异——看懂每个类的 state 存了什么、谁在改它、改到什么值时唤醒谁,这一族就只剩一张表。下一步建议挑一个你项目里正在用的同步工具,把它的 state 派生 getter(getCount()、availablePermits()、getReadLockCount()、getPhase() 之类)接进日志,观察一段时间内的真实取值曲线——比背文档有效得多。
本文关键词:AQS 共享模式、CyclicBarrier、StampedLock、线程状态、LockSupport、Phaser