Lesson 19 · 并发编程
并发工具全解:CountDownLatch、CyclicBarrier、Semaphore、Phaser
面试中四大并发工具的使用场景区别是什么?
在 java.util.concurrent 包中,有四个并发协调工具经常被放在一起考察:
- CountDownLatch —— 一次性倒计数器,一个或多个线程等待 N 个事件完成
- CyclicBarrier —— 可循环栅栏,N 个线程互相等待到达同步点后再一起继续
- Semaphore —— 信号量,控制同时访问某个资源的线程数量
- Phaser —— 灵活的阶段器,支持动态注册和阶段式推进,是前两者的超集
大多数候选人只能说出"CountDownLatch 是一次性的,CyclicBarrier 可复用"这一句——这恰恰是面试官最爱挖坑的地方。本文将从使用场景、代码示例、底层实现三个维度,彻底厘清四大工具的本质区别。
等待"事件"完成 → CountDownLatch | 等待"线程"汇合 → CyclicBarrier | 限制"并发数" → Semaphore | 动态+多阶段 → Phaser
CountDownLatch:一次性倒计数器
CountDownLatch 的语义极其简单:一个计数器,初始值为 N,每次调用 countDown() 减 1,调用 await() 的线程会阻塞直到计数器归零。计数一旦到达 0,不可重置——这是一次性工具。
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(() -> {
initService(idx); // 模拟服务初始化
latch.countDown(); // 完成一个,计数减 1
System.out.println("服务" + idx + " 初始化完毕");
}, "init-thread-" + i).start();
}
latch.await(); // 主线程阻塞,直到 5 个服务全部初始化
System.out.println("所有服务就绪,系统启动完成!");
}
}
CountDownLatch latch = new CountDownLatch(10);
AtomicInteger total = new AtomicInteger(0);
for (int i = 0; i < 10; i++) {
executor.submit(() -> {
total.addAndGet(fetchData()); // 每个线程拉取数据并累加
latch.countDown();
});
}
latch.await(); // 等待全部完成
System.out.println("汇总结果: " + total.get());
底层实现非常精炼——基于 AQS 共享模式:
private static final class Sync extends AbstractQueuedSynchronizer {
Sync(int count) { setState(count); } // AQS state = 计数值
protected int tryAcquireShared(int acquires) {
return (getState() == 0) ? 1 : -1; // state==0 获取成功,否则阻塞
}
protected boolean tryReleaseShared(int releases) {
for (;;) {
int c = getState();
if (c == 0) return false;
int nextc = c - 1;
if (compareAndSetState(c, nextc))
return nextc == 0; // 减到 0 才释放所有 await 线程
}
}
}
多个线程可能同时调用 countDown(),CAS 保证在高并发下每次减 1 的原子性。tryReleaseShared 返回 true 的条件是 nextc == 0,只有最后一个 countDown 的线程才会触发 AQS 的 doReleaseShared(),一次性唤醒所有 await() 阻塞的线程。
state = count,countDown() 通过 CAS 递减,await() 通过 AQS 共享模式阻塞直到 state==0。一次性使用,计数归零后无法重置。典型场景:等待多个服务初始化、等待多个子任务完成后汇总。
CyclicBarrier:可循环复用的栅栏
CyclicBarrier 的语义:N 个线程互相等待,全部到达屏障点后一起继续执行。与 CountDownLatch 的关键差异在于——它是"线程等线程"而非"线程等事件",并且可以反复使用。
public class MatrixCalculator {
static final int THREADS = 4;
// 所有线程到齐后,执行一次 barrierAction(如打印日志)
static CyclicBarrier barrier = new CyclicBarrier(THREADS, () ->
System.out.println("第 " + round + " 轮计算完毕,进入下一轮"));
static int round = 0;
public static void main(String[] args) {
for (int i = 0; i < THREADS; i++) {
new Thread(() -> {
for (int r = 0; r < 3; r++) {
computePartition(); // 计算自己负责的分块
try {
barrier.await(); // 等待其他线程完成本轮
} catch (Exception e) { return; }
// ★ 屏障打开后,所有线程同时进入下一轮
}
}).start();
}
}
}
底层实现与 CountDownLatch 完全不同——CyclicBarrier 没有使用 AQS,而是基于 ReentrantLock + Condition:
public class CyclicBarrier {
private final ReentrantLock lock = new ReentrantLock();
private final Condition trip = lock.newCondition();
private final int parties; // 参与者数量(不可变)
private int count; // 已到达的线程数
private int generation = 0; // 代数:每次 reset 递增,防止旧 await 误唤醒
private final Runnable barrierCommand; // 屏障动作(可选)
public int await() throws ... {
lock.lock();
try {
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(); }
}
}
因为 CyclicBarrier 需要"可重置"——当所有线程到齐后,必须把 count 恢复为 parties、generation 递增,然后开始新一轮。AQS 的 state 是单向递减模型,不支持这种"到 0 后自动恢复"的语义。Lock + Condition 天然支持这种循环模式。
| 维度 | CountDownLatch | CyclicBarrier |
|---|---|---|
| 等待关系 | 一个/多个线程等待 N 个事件完成 | N 个线程互相等待到达同步点 |
| 可复用 | 不可以(一次性) | 可以(自动重置 / reset()) |
| 底层实现 | AQS 共享模式 | ReentrantLock + Condition |
| 屏障动作 | 无 | 支持(barrierAction) |
| 异常传播 | 无(各线程独立) | 一个线程异常 → BrokenBarrierException |
"CountDownLatch 是一次性的,await() 阻塞直到 count 减为 0,底层是 AQS 共享模式。CyclicBarrier 可循环复用,N 个线程全部 await() 后才一起放行,底层是 Lock + Condition。CountDownLatch 是'等事件',CyclicBarrier 是'等线程'。"
Semaphore:基于许可的并发限流器
Semaphore 维护一组许可(permits)。acquire() 获取一个许可(没有则阻塞),release() 归还一个许可。本质是一个共享计数器,控制同时访问资源的线程数。
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(); // ★ 必须归还,否则许可泄漏!
}
}
}
// 令牌桶简化版:Semaphore + 定时补充
Semaphore rateLimiter = new Semaphore(100);
// 定时任务每秒补满许可
scheduler.scheduleAtFixedRate(() -> {
int available = rateLimiter.availablePermits();
if (available < 100)
rateLimiter.release(100 - available); // 补充到 100
}, 1, 1, TimeUnit.SECONDS);
// 请求处理
if (rateLimiter.tryAcquire()) { // 非阻塞尝试
handleRequest();
} else {
rejectWith429(); // 限流拒绝
}
Semaphore 同时支持公平模式和非公平模式,构造时通过参数选择:
// 非公平(默认):允许"插队",吞吐量高
Semaphore unfair = new Semaphore(10);
// 公平:严格按排队顺序获取许可
Semaphore fair = new Semaphore(10, true);
底层实现同样基于 AQS 共享模式,但 state 语义与 CountDownLatch 不同:
protected int tryAcquireShared(int acquires) {
for (;;) {
int available = getState();
int remaining = available - acquires;
// remaining < 0 → 许可不足,获取失败 → 入队 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 总是成功,唤醒等待线程
}
}
能,但不推荐。许可数为 1 的 Semaphore 确实是互斥的,但它没有"持有者"概念——任何线程都能 release(),不要求是 acquire 的那个线程。ReentrantLock 要求同一线程 lock/unlock。此外 Semaphore 不支持可重入。
AQS state = 可用许可数。acquire() CAS 递减 state,不够则阻塞;release() CAS 递增 state 并唤醒等待线程。支持公平/非公平模式。典型场景:连接池限流、API 限流、资源池管理。
Phaser:CountDownLatch + CyclicBarrier 的超集
Phaser 是 JDK 7 引入的灵活同步工具,解决了前两者的三个局限:参与者数量固定、只能单次或固定循环、不支持层次化。Phaser 支持动态注册/注销、阶段式推进、树形结构。
public class IterativeSolver {
static final int MAX_ROUNDS = 100;
public static void main(String[] args) {
Phaser phaser = new Phaser(1); // 主线程先注册
for (int i = 0; i < 4; i++) {
phaser.register(); // 动态注册工作线程
new Thread(() -> {
for (int r = 0; r < MAX_ROUNDS; r++) {
compute();
phaser.arriveAndAwaitAdvance(); // 到达并等待本轮所有线程
// 某些线程可能提前 deregister()
}
phaser.arriveAndDeregister(); // 完成后注销
}).start();
}
phaser.arriveAndDeregister(); // 主线程注销
}
}
Phaser 的核心 API 一览:
| 方法 | 作用 |
|---|---|
register() | 动态注册一个参与者,parties 数 +1 |
bulkRegister(n) | 批量注册 n 个参与者 |
arrive() | 到达但不等待,直接继续执行 |
arriveAndAwaitAdvance() | 到达并等待本阶段所有线程到齐(最常用) |
arriveAndDeregister() | 到达并注销自己,后续阶段不再参与 |
awaitAdvance(phase) | 等待指定阶段完成(自己不参与计数) |
Phaser 的同步状态需要同时编码"阶段号"和"到达计数"两个维度,且需要支持原子性的阶段推进。AQS 的单一 state 无法满足。Phaser 内部使用 volatile long state(64 位),高 32 位存 phase,低 32 位分两段存 parties 和 arrived 计数,通过 CAS 操作实现原子更新。
当 CountDownLatch 或 CyclicBarrier 无法满足需求时考虑 Phaser:参与者数量在运行时变化、需要多个同步阶段、需要层次化结构(父子 Phaser 降低竞争)。代价是 API 更复杂,简单场景用前两者即可。
横向对比:一张表看清四大工具
| 特性 | CountDownLatch | CyclicBarrier | Semaphore | Phaser |
|---|---|---|---|---|
| 可复用 | 不可以 | 可以 | 可以 | 可以 |
| 参与者数量 | 构造时固定 | 构造时固定 | 无"参与者"概念 | 动态注册/注销 |
| 屏障动作 | 无 | 支持 barrierAction | 无 | 支持 onAdvance() |
| 底层实现 | AQS(共享模式) | Lock + Condition | AQS(共享模式) | 自定义 volatile long + CAS |
| 公平模式 | 无 | 无 | 支持 | 无 |
| 超时支持 | await(timeout) | await(timeout) | tryAcquire(timeout) | awaitAdvanceInterruptibly(timeout) |
| 中断响应 | 支持 | 支持 | 支持 | 支持 |
| 层次化 | 不支持 | 不支持 | 不支持 | 支持(父子树形) |
根据场景选择合适的工具:
| 业务场景 | 推荐工具 | 原因 |
|---|---|---|
| 主线程等待 N 个子任务完成 | CountDownLatch | 一次性等待,语义最清晰 |
| 多线程并行计算 + 多轮同步 | CyclicBarrier | 可循环复用,支持 barrierAction |
| 限制并发访问数(连接池/限流) | Semaphore | 唯一专为"限流"设计的工具 |
| 迭代算法,线程数每轮可能变化 | Phaser | 支持动态注册,阶段式推进 |
| 多个子系统启动完毕后才能服务 | CountDownLatch | 经典场景,一次性等待 |
| 分布式仿真/游戏回合制 | Phaser | 多阶段 + 动态参与者 |
需要限流? → Semaphore
等待"事件"(一次性)? → CountDownLatch
等待"线程"汇合(可循环)? → CyclicBarrier
参与者动态变化 / 多阶段? → Phaser
决策树与高频面试陷阱
下面汇总面试中最常见的陷阱问题及标准回答:
不能。CountDownLatch 只暴露了 countDown() 方法(递减),没有 increment。如果需要"计数可以增加也可以减少"的场景,应该用 Semaphore 或自定义同步器。getCount() 可以查看当前计数值,但只用于监控/调试。
如果一个线程在 await() 期间抛出异常或被中断,CyclicBarrier 会进入 broken 状态,所有其他正在等待的线程都会收到 BrokenBarrierException。这是与 CountDownLatch 的重要区别——CountDownLatch 中各线程互相独立,一个线程异常不影响其他线程。
不需要。任何线程都可以调用 release() 归还许可。这是 Semaphore 和 Lock 的本质区别——Lock 要求同一个线程 lock/unlock(ReentrantLock 有 ownerThread 校验),Semaphore 没有"持有者"概念。但也正因如此,Semaphore 不能替代互斥锁。
功能上部分可以,但语义不同。new CyclicBarrier(N+1) 配合 N 个工作线程 + 1 个主线程可以模拟 CountDownLatch 的"等待 N 个事件"。但 CountDownLatch 更轻量(无锁实现),且语义更精确——countDown 线程完成后可以继续做别的事,而 CyclicBarrier 的 await 线程必须等到所有人到齐才能继续。
当参与者数量很大时,Phaser 支持层次化(树形)结构:多个子 Phaser 注册到一个父 Phaser 上,每个子 Phaser 独立管理一组线程。这避免了所有线程竞争同一把锁(CyclicBarrier 的 ReentrantLock),显著降低竞争。对于 20+ 线程的大规模并行计算,Phaser 的性能优势明显。
面试回答模板
Q:请比较 CountDownLatch、CyclicBarrier、Semaphore、Phaser 的区别和适用场景。
"这四个工具都是 JUC 包中的并发协调器,但解决的问题各有不同。
CountDownLatch 是一次性倒计数器,基于 AQS 共享模式实现。state 初始为 N,每次 countDown() 通过 CAS 减 1,await() 阻塞到 state==0。典型场景是主线程等待多个子任务完成。不可复用。
CyclicBarrier 是可循环栅栏,基于 ReentrantLock + Condition 实现。N 个线程全部 arrive 后一起放行,支持 barrierAction。一个线程异常会导致 BrokenBarrierException。典型场景是多轮并行计算。
Semaphore 是许可计数器,基于 AQS 共享模式。acquire() 获取许可,release() 归还。支持公平/非公平模式。典型场景是连接池限流和 API 限流。
Phaser 是 JDK 7 引入的灵活替代,支持动态注册/注销和阶段式推进,底层用 volatile long 的 CAS 操作。适合参与者数量运行时变化的场景,如迭代算法。
选型时按三个问题决策:是否限流?是否一次性等事件?参与者是否动态变化?"