JAVA · Vol.II · DAY 19 · 并发编程
并发工具全解:CountDownLatch、CyclicBarrier、Semaphore、Phaser
从一场「压测事故」说起:为什么起跑要同步
先讲一个真实场景。某团队给线上服务做大促压测,压测平台上起了 10 个施压线程,每个线程 new Thread(...).start() 后立刻开始发请求。结果呢?最先启动的线程把服务打到了 100% CPU,后面的线程请求全被限流拒绝——压测数据完全失真。问题不在压测工具,而在起跑不同步:10 个线程没有在同一时刻开跑,先跑的抢占了全部资源。
修复方案简单到令人意外:用一个 CountDownLatch(1) 当「发令枪」。所有施压线程先 await() 待命,主线程准备好后 countDown() 一声令下,10 个线程同时起跑。就这一行改动,压测曲线立刻正常。
在 java.util.concurrent(JUC)包里,有四个经常被放在一起考察的并发协调工具:
- CountDownLatch(倒计数门闩)——一次性倒计数,一个或多个线程等待 N 个事件完成;
- CyclicBarrier(循环栅栏)——N 个线程互相等待,全部到达同步点后一起继续,可循环复用;
- Semaphore(信号量)——控制同时访问某资源的线程数量(限流);
- Phaser(阶段器)——支持动态注册与多阶段推进,是前两者的超集(JDK 7+)。
大多数候选人只能说出「CountDownLatch 是一次性的,CyclicBarrier 可复用」这一句——这恰恰是面试官最爱挖坑的地方。本文从使用场景、代码示例、底层实现三个维度,把四件套(外加一个交换器 Exchanger 加分项)彻底讲透。
四个工具本质是四种「等待语义」:等事件(latch,计数减到 0)/等线程到齐(barrier,计数加到 N)/限并发数(semaphore,许可借还)/等线程且分多轮(phaser,阶段推进)。选型先问自己:我在等什么?
等待「事件」完成 → CountDownLatch | 等待「线程」汇合 → CyclicBarrier | 限制「并发数」 → Semaphore | 动态 + 多阶段 → Phaser
在深入之前,先给你一套生活化类比,后面每一站都会用到:
| 工具 | 生活类比 | 等的是什么 |
|---|---|---|
| CountDownLatch | 发令枪 / 团购集合点 | N 个「事件」全部完成(一次性) |
| CyclicBarrier | 聚餐等齐所有人再开吃 | N 个「线程」全部到齐(可复用) |
| Semaphore | 停车场车位 / 餐厅桌位 | 同时占用的「名额」上限 |
| Phaser | 多轮比赛:每轮等齐再发下一轮 | 人齐 + 推进到下一阶段 |
一张图看懂四件套的「等待模型」
为什么面试官总把这四个工具放一起考?因为它们长得像——都跟「计数」有关;但它们等的对象完全不同。先把等待模型画出来:
因为「计数语义」不同:latch 的计数是单向递减、到 0 放行;barrier 的计数是每轮归位、循环使用;semaphore 的计数是可增可减的余额;phaser 的计数是分阶段、可动态调整。四种语义对应四种不同的等待与唤醒策略,混在一起只会互相拖累——这也是 JUC 把它们拆成四个类的原因。
CountDownLatch:倒计数门闩(一次性,减到 0)
CountDownLatch(倒计数门闩)的语义极其简单:计数器初始值为 N,每次调用 countDown() 减 1,调用 await() 的线程阻塞直到计数器归零。计数一旦到达 0 就永久放行、不可重置——这是一次性(one-shot)工具。门闩(latch)这个名字很形象:门闩上有 N 根插销,每完成一件事拔掉一根,全拔完门才打开。
先看最经典的用法——主线程等待多个子任务完成:
public class ServiceBootstrap {
private static final int SERVICE_COUNT = 5;
public static void main(String[] args) throws InterruptedException {
CountDownLatch latch = new CountDownLatch(SERVICE_COUNT);
for (int i = 0; i < SERVICE_COUNT; i++) {
final int idx = i;
new Thread(() -> {
try {
initService(idx); // 模拟服务初始化
} finally {
latch.countDown(); // ★ finally 里减,防止异常导致计数漏减
}
}, "init-thread-" + i).start();
}
latch.await(); // 主线程阻塞,直到 5 个服务全部初始化完
System.out.println("所有服务就绪,系统启动完成!");
}
static void initService(int idx) {
try {
Thread.sleep(300 + idx * 100L); // 模拟耗时初始化
System.out.println("服务 " + idx + " 初始化完毕");
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
}
再看不那么起眼、但生产里用得最多的用法——并发起跑器(发令枪):用 CountDownLatch(1),让所有工作线程先 await() 待命,主线程一声令下集体开跑。开篇压测事故的修复方案就是它:
public class PressureTestStarter {
public static void main(String[] args) throws InterruptedException {
final int WORKERS = 10;
CountDownLatch startGun = new CountDownLatch(1); // 发令枪:初始 1
for (int i = 0; i < WORKERS; i++) {
new Thread(() -> {
try {
startGun.await(); // 全部待命,谁都不先跑
sendRequests(); // 听到枪响,同时开始发压测请求
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}, "worker-" + i).start();
}
Thread.sleep(1000); // 主线程做准备(预热、建连接)
System.out.println("预备——跑!");
startGun.countDown(); // 计数 1 → 0,10 个线程同时起跑
}
}
几个容易被忽略的细节:
await()有两个版本:await()无限等待;await(long timeout, TimeUnit)超时版返回boolean——true表示计数确实归零,false表示超时还没归零。生产里主线程等子任务,几乎都要用超时版,避免子任务挂了主线程永远卡死。- 计数归零后:所有正在
await()的线程被一次性放行,之后新来的await()立即返回,不会再阻塞。 countDown()放finally:任务抛异常也要减,否则计数永远到不了 0,等待方永久阻塞——这是线上最常见的 latch 事故。getCount():只读当前计数,用于监控/日志,不能用来「回加」。
一次性倒计数门闩:countDown() 减 1、await() 阻塞到 0;主线程等子任务、并发起跑器两个经典场景;生产务必用超时版 await + finally 里 countDown。
CountDownLatch 底层:AQS 共享模式与 state
CountDownLatch 的实现非常精炼,核心是一个内部类 Sync,直接复用 AQS(AbstractQueuedSynchronizer,抽象队列同步器)的共享模式。AQS 的 state 在这里就是计数值:
private static final class Sync extends AbstractQueuedSynchronizer {
Sync(int count) { setState(count); } // 构造:state = 计数值 N
protected int tryAcquireShared(int acquires) {
return (getState() == 0) ? 1 : -1; // state==0 放行,否则加入等待队列
}
protected boolean tryReleaseShared(int releases) {
for (;;) { // 自旋 CAS,保证并发 countDown 原子性
int c = getState();
if (c == 0) return false; // 已经到 0,无需再释放
int nextc = c - 1;
if (compareAndSetState(c, nextc))
return nextc == 0; // ★ 只有减到 0 的那次才返回 true
}
}
}
多个线程可能同时调用 countDown()。如果直接 state--,读-改-写三步不是原子的,两个线程同时减会丢更新。CAS 自旋保证每次减 1 都原子生效。更关键的是返回值:tryReleaseShared 只有在 nextc == 0——也就是最后一个线程 countDown 时才返回 true,AQS 收到 true 后调用 doReleaseShared(),把等待队列里所有 await() 的线程一次性唤醒。这就是「最后一个完成任务的人负责开门」。
为什么 CountDownLatch 不可重置?源码注释写得很直白:"This is a one-shot phenomenon -- the count cannot be reset. If you need a version that resets the count, consider using a CyclicBarrier."(这是单次现象,计数不可重置。如果需要可重置的版本,考虑用 CyclicBarrier。)设计上刻意不提供「加回」方法——因为可重置的倒计数是另一个需求,已经有专门的工具了。
- CountDownLatch 直接复用 AQS 共享模式:state = 剩余计数;
- await() →
acquireSharedInterruptibly:state==0 立即放行,否则入队 park; - countDown() →
releaseShared:CAS 减 1,最后一个 countDown 唤醒全部等待者; - 一次性:没有回加接口,需要可重置就用 CyclicBarrier。
CyclicBarrier:循环栅栏(可复用,等齐 N 个线程)
CyclicBarrier(循环栅栏)的语义:N 个线程互相等待,全部到达屏障点后一起继续执行。栅栏(barrier)这个比喻很贴切:跑道上横着一道栅栏,必须等所有选手都到栅栏前,栅栏才升起,大家同时继续跑。
和 CountDownLatch 的三个关键差异:① 它是「线程等线程」,每个线程自己 await() 并参与计数;② 屏障打开后自动重置,可以循环使用;③ 支持可选的 barrierAction——所有线程到齐时由最后一个到达的线程执行一段动作(比如打印、汇总)。
典型场景:分阶段(多轮)并行任务,每轮结束需要对齐一次。比如并行计算里常见的「每轮迭代完必须等所有人同步」:
import java.util.concurrent.BrokenBarrierException;
import java.util.concurrent.CyclicBarrier;
import java.util.concurrent.atomic.AtomicInteger;
public class IterationSync {
static final int THREADS = 3;
static final AtomicInteger round = new AtomicInteger(1);
// 所有线程到齐后,由最后一个到达的线程执行 barrierAction
static final CyclicBarrier barrier = new CyclicBarrier(THREADS, () ->
System.out.println("── 第 " + round.getAndIncrement() + " 轮结束,所有人到齐 ──"));
public static void main(String[] args) {
for (int i = 0; i < THREADS; i++) {
final int tid = i;
new Thread(() -> {
try {
for (int r = 0; r < 3; r++) {
computePartition(tid); // 算自己负责的分块
int index = barrier.await(); // 到齐才放行;返回到达序号
// 屏障打开:所有线程同时进入下一轮
}
} catch (BrokenBarrierException e) {
System.out.println("barrier 被破坏:" + e);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}, "worker-" + i).start();
}
}
static void computePartition(int tid) {
try {
Thread.sleep(100 + tid * 50L); // 模拟分块计算,耗时各不相同
System.out.println("worker-" + tid + " 完成一块");
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
}
注意代码里 barrier.await() 的返回值:它是当前线程在本轮的到达序号——最后一个到达的线程返回 0,其余按到达顺序返回 N-1、N-2…。这个返回值偶尔有用(比如最后到达的线程负责干收尾活),但多数场景直接忽略。
barrierAction 由最后一个到达的线程在「所有人到齐」的瞬间执行,执行完才统一放行。所以它适合做「每轮汇总/打日志」这类轻量工作;千万别在里面做重活——它会拖住整轮所有线程。如果汇总很重,让 barrierAction 只负责「发信号」,重活在别的线程做。
CyclicBarrier 底层:Lock + Condition 与 broken 机制
CyclicBarrier 的底层和 CountDownLatch 完全不同——它没有直接使用 AQS,而是基于 ReentrantLock + Condition(ReentrantLock 本身是 AQS 独占模式实现,所以算「间接」地基)。为什么这么设计?因为「可重置」:一轮结束后 count 要恢复为 parties、代数要递增,AQS 的 state 单向递减模型不支持这种「归零后自动还原」的循环语义,而 Lock + Condition 天然适合。
public class CyclicBarrier {
private final ReentrantLock lock = new ReentrantLock();
private final Condition trip = lock.newCondition();
private final int parties; // 参与者总数(构造时定死)
private int count; // 还差几个到齐(从 parties 往下减)
private int generation = 0; // 代数:每轮 +1,防止旧轮线程误唤醒
private final Runnable barrierCommand; // 可选的 barrierAction
public int await() throws InterruptedException, BrokenBarrierException {
lock.lock();
try {
if (broken) throw new BrokenBarrierException();
int g = generation;
int index = --count; // 到达一个,差数减 1
if (index == 0) { // ★ 最后一个到达的线程
if (barrierCommand != null) barrierCommand.run();
nextGeneration(); // 唤醒所有等待者,count 复位,generation++
return 0;
}
while (g == generation) // 等同一代:防止上一轮残留的唤醒
trip.await(); // 在 Condition 上挂起
return index;
} finally { lock.unlock(); }
}
}
三个设计细节值得记住:
- generation(代数):每轮结束
generation++,等待线程用g == generation校验自己属于哪一轮,避免「上一轮的唤醒信号」把下一轮线程提前放走; - count 复位:
nextGeneration()里把 count 重新置回 parties,实现自动循环复用; - broken 状态:如果某个线程在 await 期间被中断、或 await 超时、或 barrierAction 抛异常,整个栅栏进入 broken(破坏)状态——所有正在等待的线程立刻收到
BrokenBarrierException,之后新来的 await 也直接抛异常。想恢复只能调reset()。
这是高频考点:一个线程异常 → 整代 broken → 所有等待线程抛 BrokenBarrierException。这与 CountDownLatch 截然不同——latch 里各线程相互独立,一个线程挂了不影响别的线程继续 countDown。所以面试官总爱用这个问题区分「你到底是真懂,还是只会背 API」。
- 基于 ReentrantLock + Condition(AQS 的间接复用),靠 generation 区分轮次实现循环;
- N 个线程全部 await() 才放行,最后一个到达的线程执行 barrierAction;
- 一个线程中断/超时 → 整代 broken,全员
BrokenBarrierException,reset()可恢复; - await() 返回到达序号,最后一个到达返回 0。
深度对比:latch「减到 0」vs barrier「加到 N」
面试官最爱的追问来了:「CountDownLatch 和 CyclicBarrier 到底什么区别?」光说「一个一次性一个可复用」不够,要从计数方向和等待语义两个层面说透。
- 计数方向相反:latch 是「事件完成数」,从 N 减到 0;barrier 是「已到齐线程数」,从 0 加到 N(实现里 count 从 parties 递减,但语义是等齐 N 个)。
- 等待者不同:latch 的
await()是「旁观者」——等事件的线程(通常是主线程)自己不参与计数;barrier 的await()是「参与者」——每个线程都await()并贡献一次计数。 - await 之后的去向不同:latch 里 countDown 的线程完成事件后可以立刻继续干别的,不用等别人;barrier 里每个线程都必须等齐所有人才能继续。
- 复用性不同:latch 计数到 0 后永久放行;barrier 每轮自动复位。
barrier:等「线程」,计数 0 → N,可复用,每个线程必须等齐所有人
① 能互相替代吗?「new CyclicBarrier(N+1) + N 个工作线程 + 1 个主线程」确实可以模拟 latch 的「等 N 个事件」——但语义不同:latch 的工作线程 countDown 后立刻继续干活,barrier 模拟时工作线程必须原地等齐所有人。反过来,latch 无法模拟 barrier 的「循环复用」。
② CountDownLatch 是「无锁实现」吗?不是。它基于 AQS,内部同样用 CAS 与 park/unpark,只是实现更薄(没有额外的 Condition、没有持有者概念)。面试说「latch 无锁」会被扣分。
Semaphore:信号量(许可借还,限并发数)
Semaphore(信号量)维护一组许可(permits):acquire() 拿一个许可(没有就阻塞),release() 还一个许可。本质是一个可增可减的共享计数器,用来控制同时访问某个资源的线程数量。类比停车场:车位总数固定,进一辆占一个位,没位了就在门口排队,出一辆放一辆进来。
第一个场景——限制数据库连接池的并发使用数(注意 release 必须在 finally):
public class DbPool {
// 最多 10 个线程同时持有连接
private static final Semaphore semaphore = new Semaphore(10);
public ResultSet query(String sql) throws InterruptedException {
semaphore.acquire(); // 没许可就阻塞排队
try {
return executeQuery(sql);
} finally {
semaphore.release(); // ★ 必须归还,否则许可泄漏 → 全阻塞
}
}
}
第二个场景——API 限流(令牌桶简化版):注意 Semaphore 是「计数信号量」不是「速率限制器」,它只能限制同时处理数,不能平滑限速;做「每秒最多 X 个」需要配合定时任务补许可(真正的令牌桶可参考 Guava RateLimiter):
Semaphore rateLimiter = new Semaphore(100);
// 定时任务每秒把许可补回 100
scheduler.scheduleAtFixedRate(() -> {
int available = rateLimiter.availablePermits();
if (available < 100)
rateLimiter.release(100 - available);
}, 1, 1, TimeUnit.SECONDS);
// 请求处理:抢不到许可立刻降级,别让请求干等
if (rateLimiter.tryAcquire()) { // 非阻塞尝试
handleRequest();
} else {
rejectWith429(); // 限流拒绝
}
Semaphore 支持公平 / 非公平两种模式,由构造参数决定:
// 非公平(默认):允许「插队」,吞吐量更高
Semaphore unfair = new Semaphore(10);
// 公平:严格按排队顺序获取许可(FIFO),适合对公平性敏感的场景
Semaphore fair = new Semaphore(10, true);
除了 acquire()/release(),还有几个变体:tryAcquire() 非阻塞抢许可;tryAcquire(timeout, unit) 限时抢(抢不到就超时返回,适合「有超时要求的限流」);acquire(n)/release(n) 批量拿还;availablePermits() 查看余量;drainPermits() 一次拿光(用于停机时排空)。
许可借还模型:acquire 减许可(可阻塞/可超时/可中断),release 加许可;连接池控制、API 限流、令牌桶简化是三大场景;release 必须放 finally;默认非公平,需要 FIFO 传 true。
Semaphore 底层与三个坑
Semaphore 同样复用 AQS 共享模式,state 的语义是「当前可用许可数」。非公平模式拿许可时先 CAS 抢一次(能抢到就不排队),抢不到再入队:
protected int tryAcquireShared(int acquires) {
for (;;) {
int available = getState();
int remaining = available - acquires;
// remaining < 0 → 许可不足,获取失败 → AQS 入队 park
// CAS 成功 → 许可扣减成功
if (remaining < 0 || compareAndSetState(available, remaining))
return remaining;
}
}
protected boolean tryReleaseShared(int releases) {
for (;;) {
int current = getState();
int next = current + releases;
if (compareAndSetState(current, next))
return true; // release 几乎总是成功,并唤醒等待线程
}
}
看懂这段源码,三个「坑」就都懂了——这三条也是面试官最爱埋的点:
能,而且这是 bug 源头。release() 没有「必须持有许可」的校验,谁都能调。常见错误:if (tryAcquire()) { 处理 } release(); —— 把 release 写在 if 外面,抢不到许可的线程也 release 一次,许可被凭空造出来,限流失效、甚至许可数超过初始值。正确写法:拿到许可才归还,且放 finally。
会。Semaphore(5) 初始 5 个许可,如果某个线程多调了一次 release(),许可变成 6。所以「许可数 = 初始值」不是不变量——用 availablePermits() 监控时别拿它当常量。
能保证互斥,但不推荐:① 没有「持有者」概念,线程 A acquire、线程 B release 也合法,锁的语义被打破;② 不可重入——同一线程 acquire 两次会把自己阻塞死(ReentrantLock 可重入);③ 非公平下抢锁顺序无保证。需要互斥锁用 synchronized / ReentrantLock。
- state = 可用许可数,acquire 减、release 加,全程 CAS;
- release 无持有者校验:tryAcquire 失败后误 release 会「凭空造许可」;
- release 可超额,许可数可能超过初始值;
- Semaphore(1) 不等于互斥锁(无重入、无持有者)。
Phaser:阶段器(动态参与者 + 多阶段,JDK 7+)
Phaser(阶段器,JDK 7 引入)是前两者的「超集」,解决三个痛点:① 参与者数量固定——Phaser 支持运行时动态注册/注销(register()/arriveAndDeregister());② 只能单次或固定循环——Phaser 按「阶段(phase)」推进,每个阶段都是一个迷你屏障;③ 支持树形结构(父子 Phaser)降低大规模竞争。
一句话理解:Phaser = 多轮的 CyclicBarrier,且每轮参赛人数可以变。适用场景:迭代算法、流水线、游戏回合制——线程数在运行中会增减的任务。
import java.util.concurrent.Phaser;
public class DynamicSolver {
static final int ROUNDS = 5;
public static void main(String[] args) {
Phaser phaser = new Phaser(1); // 主线程先占一个位(注册)
for (int i = 0; i < 4; i++) {
phaser.register(); // 每来一个工作线程,动态注册 +1
new Thread(() -> {
for (int r = 0; r < ROUNDS; r++) {
compute(); // 干活
phaser.arriveAndAwaitAdvance(); // 到达本阶段并等齐
}
phaser.arriveAndDeregister(); // 干完活:到达并注销自己(参与数 -1)
}, "worker-" + i).start();
}
while (phaser.getRegisteredParties() > 1) { // 主线程盯进度
phaser.arriveAndAwaitAdvance(); // 每轮也跟着对齐
}
System.out.println("所有轮次完成,主线程收工");
}
static void compute() {
try {
Thread.sleep(50); // 模拟计算
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
}
Phaser 的核心 API 一览:
| 方法 | 作用 |
|---|---|
register() / bulkRegister(n) | 动态注册 1 个 / n 个参与者,parties +1 / +n |
arrive() | 到达但不等待,立即继续执行 |
arriveAndAwaitAdvance() | 到达并等待本阶段所有人到齐(最常用) |
arriveAndDeregister() | 到达并注销自己,后续阶段不再参与 |
awaitAdvance(phase) | 等指定阶段结束(自己不参与计数) |
getRegisteredParties() / getPhase() | 当前参与人数 / 当前阶段号(从 0 开始) |
CyclicBarrier 的参与者在构造时定死,中途不能变;Phaser 的参与者可以随时 register() / arriveAndDeregister()。另一个差异:barrier 里每个线程到点必须 await();Phaser 里 arrive() 和 awaitAdvance() 是拆开的——你可以「只报到我到了、但不在原地等」,灵活得多,也复杂得多。
Phaser 底层:一个 long 塞下四种状态
Phaser 没有复用 AQS 的 state,而是自研了一个 volatile long state——一个 64 位变量同时编码四种信息(这是源码注释里明确写的位布局):
// 位布局(从低位到高位):
// bits 0-15 : unarrived 当前阶段还没到达的人数
// bits 16-31 : parties 总参与人数
// bits 32-62 : phase 阶段号(从 0 递增)
// bit 63 : terminated 是否已终止(1 = 终止)
private volatile long state;
// 例:参与 4 人、已到 3 人、当前第 2 阶段
// state = 0x00000002_00040003(示意)
// phase=2 parties=4 unarrived=3
阶段推进时,最后一个到达的线程用 CAS 把 unarrived 重置为 parties、phase+1,一步到位完成「开门 + 关门 + 换代」。因为参与人数、阶段号随时可能变,AQS 那种「固定 state 语义」的模型装不下,所以 Phaser 选择自己管理位运算——思想仍是 AQS 那一套(volatile + CAS),只是状态更丰富。
两个高级特性,面试提到就是加分项:
onAdvance(phase, registeredParties):每阶段结束时回调,默认实现在「注册人数归 0」时返回true让 Phaser 终止(之后arrive会抛IllegalStateException)。重写它可实现「跑满 N 阶段自动终止」;- 树形结构:
new Phaser(parentPhaser)创建子 Phaser,大量参与者分到多个子 Phaser 上各自计数,子 Phaser 只向父 Phaser 汇报,把「万人抢一把锁」的竞争摊开——20+ 线程的大规模并行里比 CyclicBarrier 快。
参与者数量运行时会变、或者需要多阶段推进 → 用 Phaser;只是「等一次」或「固定轮次对齐」→ 用 latch / barrier,别为复杂度买单。Phaser 是四件套里最强大也最复杂的,简单场景用它属于过度设计。
Exchanger:两线程交换数据(加分项)
Exchanger(交换器,JDK 1.5+)是 JUC 里存在感最低的工具之一,但面试偶尔会带一句:两个线程在同一个点交换各自的数据。exchange(x) 把 x 交给对方,同时拿到对方的数据;如果没有对手,就阻塞等待(可换超时版 exchange(x, timeout, unit))。
经典场景是双缓冲:生产者线程填满 buffer A,消费者线程在 buffer B 上处理数据,两个线程通过 Exchanger 互换 buffer,避免加锁拷贝:
import java.util.concurrent.Exchanger;
import java.util.ArrayList;
import java.util.List;
public class BufferSwap {
public static void main(String[] args) {
Exchanger<List<Integer>> exchanger = new Exchanger<>();
new Thread(() -> { // 生产者
List<Integer> buf = new ArrayList<>();
try {
for (int i = 0; i < 10; i++) {
buf.add(i); // 攒一批数据
if (buf.size() == 5) {
buf = exchanger.exchange(buf); // 满 5 个就交换,换回空 buffer
}
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}, "producer").start();
new Thread(() -> { // 消费者
List<Integer> buf = new ArrayList<>();
try {
for (int i = 0; i < 2; i++) {
buf = exchanger.exchange(buf); // 换回装满的 buffer
System.out.println("消费者处理: " + buf); // 处理
buf.clear(); // 清空再换回去
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}, "consumer").start();
}
}
几个使用注意:① Exchanger 只支持两方配对,三方以上得拆成多组;② 不配对会一直阻塞,务必用超时版或保证双方必然到达;③ 底层是 CAS + 自旋(slot 节点),线程多时性能一般——它天生是为「两线程高频互换」设计的。
两线程在汇合点交换数据,典型场景是双缓冲(避免拷贝);只支持两方、会阻塞、要防死等;面试提一句「双缓冲 + 超时保护」就够加分了。
五大工具对比总表与选型决策树
把五件套拉一张总表(Exchanger 一起比,虽然它最特殊):
| 特性 | CountDownLatch | CyclicBarrier | Semaphore | Phaser | Exchanger |
|---|---|---|---|---|---|
| 等什么 | N 个事件完成 | N 个线程到齐 | 拿许可 | 阶段内人齐 | 对方线程到达 |
| 可复用 | 不可以 | 可以(自动/reset) | 可以 | 可以 | 可以(反复换) |
| 参与者数量 | 构造时固定 | 构造时固定 | 无参与者概念 | 动态注册/注销 | 固定两方 |
| 参与者是否等待 | countDown 不等别人 | 全部 await 互等 | acquire 等许可 | 可到齐可不等 | 双方互等 |
| 屏障动作 | 无 | barrierAction | 无 | onAdvance() | 无 |
| 底层实现 | AQS 共享 | Lock + Condition | AQS 共享 | volatile long + CAS | CAS + 自旋 |
| 公平模式 | 无 | 无 | 支持 | 无 | 无 |
| 超时支持 | await(timeout) | await(timeout) | tryAcquire(timeout) | awaitAdvance 超时版 | exchange(x, timeout) |
| 中断响应 | 支持 | 支持(会 broken) | 支持 | 支持 | 支持 |
| 层次化 | 不支持 | 不支持 | 不支持 | 支持(树形) | 不支持 |
| 典型场景 | 主线程等子任务、起跑 | 多轮并行计算 | 连接池、限流 | 动态流水线、回合制 | 双缓冲 |
选型时走一遍决策树,三个问题锁定答案:
需要限流? → Semaphore
等待「事件」(一次性)? → CountDownLatch
等待「线程」汇合(可循环)? → CyclicBarrier
参与者动态变化 / 多阶段? → Phaser
两线程互换数据? → Exchanger
底层统一视角:一切皆 AQS 的 state
把四件套的底层放在一起看,会发现一个统一规律:几乎所有并发同步器都是「对 AQS 的 state 做文章」。state 是一个 int(或 long),每种工具给它定义不同的语义:
| 同步器 | state 语义 | 增减方向 | 与 AQS 的关系 |
|---|---|---|---|
| CountDownLatch | 剩余事件数 | 减到 0 放行 | 直接复用(共享模式) |
| Semaphore | 可用许可数 | 可加可减 | 直接复用(共享模式) |
| ReentrantLock | 重入次数(0 = 无锁) | 加锁 +1 / 解锁 -1 | 直接复用(独占模式) |
| CyclicBarrier | (无独立 state) | count 从 N 递减 | 间接:经 ReentrantLock + Condition |
| Phaser | long 位打包:unarrived / parties / phase / terminated | 阶段推进时重置 | 自研 volatile + CAS(思想同源) |
| Exchanger | (无 state) | slot CAS | 无锁实现(CAS + 自旋) |
再往深一层,AQS 的设计是模板方法模式:排队、阻塞、唤醒这些「通用骨架」由 AQS 写死,而「怎么算获取成功」留给子类实现——CountDownLatch 实现 tryAcquireShared(state==0 才算成功),Semaphore 实现 tryAcquireShared(剩余许可 ≥ 0 才算成功)。你学会一个,其它都是换皮。
遇到并发协调需求,别急着翻 API,先想清楚:我需要的计数是「减到 0」(latch)、「凑到 N」(barrier)、「余额可借还」(semaphore)还是「分阶段可变」(phaser)?想清楚语义,工具就自己冒出来了。这也是面试官想听的思考路径——不是背 API,是理解 state。
生产实战:起跑、汇总、限流、流水线
把四件套放进一个真实的压测平台里,看看它们各司其职:
- 并发起跑(CountDownLatch(1)):10 个施压线程 await 待命,主线程 countDown 集体开跑——开篇事故的修复;
- 批量任务汇总(CountDownLatch(N)):主线程提交 100 个压测任务到线程池,每个任务完成 countDown,主线程 await(timeout) 等全部完成再出报表——注意用超时版,防止某个任务卡死导致报表永远出不来;
- 限流保护(Semaphore):压测目标服务有配额,Semaphore 限制同时打进去的请求数,避免把服务真打挂;
- 多轮流水线(CyclicBarrier / Phaser):压测分多轮(预热轮 → 正式轮 → 爬坡轮),每轮结束对齐一次再进入下一轮。
围绕这四个用法,汇总一份「踩坑清单」——每一行都是线上真实教训:
| 坑 | 后果 | 正确做法 |
|---|---|---|
| countDown() 没放 finally,任务抛异常 | 计数到不了 0,await 永久阻塞 | countDown 放 finally,或用 try-with-resources 风格封装 |
| await() 用无超时版 | 子任务挂了主线程跟着挂 | 生产一律 await(timeout) 并处理 false |
| latch 想复用 | 计数归零后失效 | 每轮 new 一个,或改用 CyclicBarrier / Phaser |
| barrier 中一个线程超时 | 整代 broken,全员 BrokenBarrierException | catch BrokenBarrierException 后 reset() 或整体重试 |
| tryAcquire 失败后仍 release | 许可凭空增加,限流失效 | 只在拿到许可的分支里 release,且放 finally |
| Semaphore 忘了 finally release | 许可泄漏,最终全部阻塞 | acquire / release 严格 try-finally 配对 |
| 线程池 + latch 任务被拒 | countDown 次数不够,主线程干等 | 拒绝策略里也 countDown,或用 CompletableFuture 编排 |
| Phaser 线程忘了 arriveAndDeregister | 注册数虚高,阶段永远推不进 | 工作线程收尾务必注销,配合 getRegisteredParties 监控 |
- 「等一次」用 CountDownLatch,「等齐循环」用 CyclicBarrier,「限并发」用 Semaphore,「动态多阶段」才上 Phaser;
- 所有等待一律给超时;所有 release / countDown 一律 try-finally;
- latch / barrier 和线程池配合时,把「任务没执行(被拒 / 异常)」也算作完成事件,否则计数对不上;
- 监控三件套:
getCount()/availablePermits()/getRegisteredParties()打进日志,出事能定位。
面试 30 秒总结与高频陷阱
先把「面试 30 秒总结」背下来,然后看陷阱题——总结负责拿分,陷阱题负责防扣分:
"JUC 的并发协调工具本质是四种计数语义:CountDownLatch 等事件,计数从 N 减到 0,一次性,基于 AQS 共享模式,主线程等子任务、并发起跑器;CyclicBarrier 等线程,N 个线程全部 await 才放行,自动重置可复用,基于 ReentrantLock + Condition,一个线程异常全员 BrokenBarrierException;Semaphore 限并发,许可借还、可公平可超时,基于 AQS 共享模式,连接池和 API 限流;Phaser 是超集,动态注册 + 多阶段,JDK 7+,参与者会变的场景才用它。选型先问自己在等什么。"
不能。只暴露 countDown()(减),没有增加接口,源码注释明确这是 one-shot。需要「可增可减的计数」用 Semaphore;需要「可重置的等待」用 CyclicBarrier。getCount() 只读,用于监控。
该线程在 await 中被中断/超时 → 栅栏进入 broken 状态 → 所有正在等待的线程收到 BrokenBarrierException,后续 await 也直接抛。CountDownLatch 里各线程互相独立,没有这个连锁反应——这是两者最本质的行为差异之一。
不需要。release 没有持有者校验,任何线程都能归还(甚至超额归还)。这正是 Semaphore 不能替代互斥锁的原因——ReentrantLock 有 owner 校验、可重入,Semaphore 两者都没有。
latch 的 await(timeout) 返回 boolean(是否归零),且 countDown 的线程不等待;barrier 的 await() 返回到达序号(最后到达为 0),await(timeout) 超时会抛 TimeoutException 并把栅栏打坏。同样是「等待」,返回值与失败模型完全不同。
① 参与者数量运行时会变(register / deregister);② 需要多阶段推进;③ 参与者很多(20+)时用树形 Phaser 把竞争分摊,避免所有人抢 CyclicBarrier 那一把 ReentrantLock。
latch 减到 0 放行、barrier 凑到 N 放行、semaphore 借还许可、phaser 动态多阶段——底层全是 state 上的 CAS,选型先想计数语义。
这一篇你掌握了什么
核心知识点回顾
- CountDownLatch(倒计数门闩):一次性,计数从 N 减到 0;await() 阻塞、countDown() 减 1(放 finally);主线程等子任务 + 并发起跑器(latch(1) 当发令枪);底层直接复用 AQS 共享模式,state = 剩余计数,最后一个 countDown 唤醒全部;
- CyclicBarrier(循环栅栏):可复用,N 个线程全部 await() 才放行,支持 barrierAction;底层是 ReentrantLock + Condition(AQS 间接复用),generation 分轮次;一个线程中断/超时 → 全员 BrokenBarrierException;
- Semaphore(信号量):许可借还,acquire 减 / release 加(都必须 try-finally);连接池、API 限流、令牌桶简化;公平/非公平;三个坑——tryAcquire 失败误 release 凭空造许可、release 可超额、Semaphore(1) 不是互斥锁;
- Phaser(阶段器,JDK 7+):动态注册/注销 + 多阶段推进;底层 volatile long state 位打包(unarrived / parties / phase / terminated);onAdvance 终止、树形结构降竞争;简单场景别用;
- Exchanger(交换器):两线程汇合点交换数据,双缓冲经典场景,只支持两方;
- 统一底层:四件套都是「对计数做文章」——CountDownLatch / Semaphore 直接复用 AQS state,CyclicBarrier 经 ReentrantLock,Phaser 自研位打包 + CAS;选型先想「我在等什么」。对比关键:latch 减到 0(一次性、等事件)vs barrier 加到 N(可复用、等线程)。
一句话总结
四个工具 = 四种等待语义:等事件(latch)/等线程(barrier)/限并发(semaphore)/动态多阶段(phaser),底层都建立在 AQS 的 state 之上。面试先答语义、再答底层、最后补场景,30 秒足够拿满分。
下一篇预告:JMM 与 happens-before——为什么 volatile 能保证可见性、为什么 i++ 不安全,从内存模型层面把并发的地基打牢。
Comments · 评论