← 返回博客
2026-09-24 08:00:02

AQS 共享模式只服务两件:CountDownLatch 和 Semaphore

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:

同步器家族:AQS 共享模式只服务 CountDownLatch 和 Semaphore

图里最上面那个 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,比翻文档可靠——这个类是包级私有的,只能反射拿到。

/img/rrwl-state-bits.svg

上面这张图是本文最需要注意版本差异的地方。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 布局和乐观读时序看清楚:

StampedLock 的 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 是什么同步基础一次性/可复用适合的场景
CountDownLatchAQS int state = 剩余计数AQS 共享模式一次性,归零即废一个线程等 N 个任务齐活
SemaphoreAQS int state = 可用许可AQS 共享模式可复用限流、连接池、一许可当互斥锁
CyclicBarrier自己的 int count + GenerationReentrantLock + Condition可复用,但会破损N 个线程多轮齐步走
Phaser自己的 volatile long stateVarHandle CAS + QNode可复用,可动态增删多阶段 + 参与方数量会变
Exchanger自己的 Slot[] + volatile int boundVarHandle CAS可复用两条流水线对调缓冲,一次配一对
ReentrantReadWriteLockAQLS 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。

回到面试速答,把这段最容易答错的四条重新写一遍:

核心收获就一句:这些工具类的差异不是 API 风格差异,而是 state 语义的差异——看懂每个类的 state 存了什么、谁在改它、改到什么值时唤醒谁,这一族就只剩一张表。下一步建议挑一个你项目里正在用的同步工具,把它的 state 派生 getter(getCount()、availablePermits()、getReadLockCount()、getPhase() 之类)接进日志,观察一段时间内的真实取值曲线——比背文档有效得多。

本文关键词:AQS 共享模式、CyclicBarrier、StampedLock、线程状态、LockSupport、Phaser