Code Review 时发现同事用 CyclicBarrier 写了个一次性等待
9 月中的一次 code review,我看到同事写了这么一段:一个接口要并行查三个数据源然后汇总,他用 CyclicBarrier 做同步。功能是对的,但读起来很别扭——CyclicBarrier 那套 await() 加 BrokenBarrierException 的写法,用在一个只同步一次的场景上,纯属给自己找麻烦。
我让他换成 CountDownLatch,代码短了三分之一。后来在组里分享了一次这三个工具的差别,顺便把自己以前用错的地方也复盘了。
CountDownLatch:等人齐,一次性
语义是"一个(或多个)线程等着,直到其他 N 个线程干完活"。内部用 AQS 的共享模式,state 就是计数值:
public class CountDownLatch {
private final Sync sync;
public void await() throws InterruptedException {
sync.acquireSharedInterruptibly(1);
}
public void countDown() {
sync.releaseShared(1);
}
private static final class Sync extends AbstractQueuedSynchronizer {
Sync(int count) {
setState(count); // state = 计数
}
protected int tryAcquireShared(int acquires) {
return (getState() == 0) ? 1 : -1; // 0 才算获取成功
}
protected boolean tryReleaseShared(int releases) {
for (;;) { // CAS 自旋减一
int c = getState();
if (c == 0) return false;
int nextc = c - 1;
if (compareAndSetState(c, nextc))
return nextc == 0;
}
}
}
}
就这点代码。await 是"等 state 变成 0",countDown 是"state 减 1"。减到 0 时唤醒所有等待者。
关键限制:计数只能减不能加,没有 reset 方法。所以它是次性的。
典型用法,也是最常见的坑:countDown() 一定要放 finally。
public OrderDetailVO buildDetail(Long orderId) throws InterruptedException {
CountDownLatch latch = new CountDownLatch(3);
AtomicReference<Order> orderRef = new AtomicReference<>();
AtomicReference<List<OrderItem>> itemsRef = new AtomicReference<>();
AtomicReference<User> userRef = new AtomicReference<>();
executor.submit(() -> {
try { orderRef.set(orderMapper.selectById(orderId)); }
finally { latch.countDown(); } // 放 finally,查询抛异常也能减计数
});
executor.submit(() -> {
try { itemsRef.set(itemMapper.listByOrderId(orderId)); }
finally { latch.countDown(); }
});
executor.submit(() -> {
try { userRef.set(userService.getByOrderId(orderId)); }
finally { latch.countDown(); }
});
// 等 2 秒,超时就不再等了,用已有的数据兜底返回
if (!latch.await(2, TimeUnit.SECONDS)) {
log.warn("buildDetail timeout, orderId={}, remaining={}",
orderId, latch.getCount());
}
return assemble(orderRef.get(), itemsRef.get(), userRef.get());
}
两个细节:
countDown()放finally。我第一次写的时候放在方法体最后,结果某个查询抛异常,计数少减 1,主线程永远卡在await上。Tomcat 线程就这么被吃掉了,攒到 200 个接口全挂。await()带超时。这个接口 P99 之前是 3.2 秒,改成并行加 2 秒超时后降到 980 毫秒。超时后不抛异常,用已有数据降级返回,用户看得到部分信息,比白屏强。
CyclicBarrier:互相等,可复用
语义不一样:它是一组线程互相等待,所有线程都到齐了才一起往下走。名字里的 Barrier 就是"栅栏",线程跑到栅栏前停下,等其他人。
public class CyclicBarrier {
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 static class Generation {
boolean broken = false; // 栅栏是否被打破
}
}
注意它不是基于 AQS 的,用的是 ReentrantLock + Condition。它的可复用靠的是 Generation 这个内部类:一轮结束后 nextGeneration() 建一个新的 Generation,重置计数。
private int dowait(boolean timed, long nanos) throws ... {
final ReentrantLock lock = this.lock;
lock.lock();
try {
final Generation g = generation;
if (g.broken) throw new BrokenBarrierException();
int index = --count;
if (index == 0) { // 最后一个到的线程
boolean ranAction = false;
try {
final Runnable command = barrierCommand;
if (command != null)
command.run(); // 执行 barrierAction
ranAction = true;
nextGeneration(); // 开启新一轮,唤醒所有人
return 0;
} finally {
if (!ranAction)
breakBarrier(); // action 抛异常,栅栏破掉
}
}
for (;;) {
try {
if (!timed)
trip.await();
else if (nanos > 0L)
nanos = trip.awaitNanos(nanos);
} catch (InterruptedException ie) {
if (g == generation && !g.broken) {
breakBarrier(); // 被中断也会打破栅栏
throw ie;
} else {
Thread.currentThread().interrupt();
}
}
if (g.broken) throw new BrokenBarrierException();
if (g != generation) return index; // 换代了,说明到齐了
if (timed && nanos <= 0L) {
breakBarrier();
throw new TimeoutException();
}
}
} finally {
lock.unlock();
}
}
那段 breakBarrier() 的逻辑是我以前没注意的:任何一个等待线程被中断或超时,整个栅栏就"破"了,其他所有正在等待的线程都会收到 BrokenBarrierException。这是它和 CountDownLatch 很大的不同——CountDownLatch 里一个线程出问题,其他线程照常走自己的。
它的典型场景是分批处理数据。我们有个对账任务,4 个线程各处理一个分片,每处理完一批(比如 1 万条)就在栅栏处汇合,汇总一次中间结果,然后进入下一批:
public void reconcile(List<List<Bill>> shards) {
int workerNum = shards.size();
AtomicLong totalMatched = new AtomicLong();
CyclicBarrier barrier = new CyclicBarrier(workerNum, () -> {
// 所有分片都处理完当前批次后,这个 Runnable 会被最后一个到的线程执行
log.info("batch done, matched={}", totalMatched.get());
});
for (int i = 0; i < workerNum; i++) {
final List<Bill> shard = shards.get(i);
executor.submit(() -> {
for (int batch = 0; batch < shard.size(); batch += BATCH_SIZE) {
List<Bill> sub = shard.subList(batch,
Math.min(batch + BATCH_SIZE, shard.size()));
totalMatched.addAndDo(match(sub));
try {
barrier.await(); // 等其他人处理完这一批
} catch (BrokenBarrierException e) {
log.error("barrier broken, abort", e);
return; // 有同伴挂了,自己也没法继续
}
}
});
}
}
这个场景换成 CountDownLatch 就写不了——因为需要同步 200 多批,CountDownLatch 是一次性的,得创建 200 多个。
Semaphore:限流,不是锁
前两个是"同步",Semaphore 是"限流"。它维护一组许可,acquire() 拿一个,release() 还一个,拿不到就阻塞。
public class Semaphore {
private final Sync sync;
protected int tryAcquireShared(int acquires) {
for (;;) {
int available = getState();
int remaining = available - acquires;
if (remaining < 0 || compareAndSetState(available, remaining))
return remaining;
}
}
}
注意 Semaphore 是个计数信号量,许可数可以是任意正整数。设为 1 时它就退化成一个"共享锁",但和真正的锁有本质区别:Semaphore 不要求 release 的线程是 acquire 的线程。任何线程都可以 release(),凭空增加许可数。这个特性有时是 bug 的来源,有时是特性(比如用线程 A 申请、线程 B 释放来传递控制权)。
我们的用法是给下游接口限流。有个第三方物流查询接口,对方明确说并发不能超过 20,超了就限流拒绝:
@Component
public class LogisticsClient {
// 20 个许可,公平模式(先到先得,避免饥饿)
private final Semaphore semaphore = new Semaphore(20, true);
public TrackInfo query(String trackingNo) {
boolean acquired = false;
try {
// 最多等 500 毫秒,等不到就降级
acquired = semaphore.tryAcquire(500, TimeUnit.MILLISECONDS);
if (!acquired) {
log.warn("logistics query throttled, no={}", trackingNo);
return TrackInfo.unknown();
}
return httpGet(trackingNo);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
return TrackInfo.unknown();
} finally {
if (acquired) { // 一定要判断,不然凭空多发许可
semaphore.release();
}
}
}
}
那个 if (acquired) 的判断是必须的。我第一版写成了无条件 release(),结果:没拿到许可的请求也调用了 release(),凭空多出一个许可,许可总数从 20 一路涨到 40、60,限流形同虚设。压测时下游 QPS 冲到 180,把对方打挂了,被对方的技术在群里 @。这个 bug 我记了很久。
公平模式 new Semaphore(20, true) 的代价是吞吐低。我们测过,20 许可、100 并发抢,非公平模式 12.4 万 ops/s,公平模式 3.1 万 ops/s,差 4 倍。但限流场景要的是"不饿死",不是"快",所以用公平。
三个对比
| CountDownLatch | CyclicBarrier | Semaphore | |
|---|---|---|---|
| 语义 | 等人齐 | 互相等 | 限流 |
| 底层 | AQS 共享模式 | ReentrantLock + Condition | AQS 共享模式 |
| 可复用 | 否 | 是 | 是 |
| 计数方向 | 递减到 0 | 递减到 0 后重置 | 可增可减 |
| 触发动作 | 无 | barrierAction(最后一个线程执行) | 无 |
| 异常传播 | 无 | BrokenBarrierException 传染给所有人 | 无 |
| 释放者限制 | 任意线程 countDown | 必须参与者自己 await | 任意线程 release |
一句话选型:等别人干活 → CountDownLatch;多轮同步 → CyclicBarrier;控制并发数 → Semaphore。
常见误用清单
- 用 CyclicBarrier 做一次性等待。就是开头那个 case。能用 CountDownLatch 就别用 CyclicBarrier,后者的异常处理麻烦得多。
countDown()不放 finally。异常路径漏掉计数,等待线程永远卡死。await()不带超时。线上出问题时,所有请求线程堆在await上,整个服务不可用。我现在一律await(2, TimeUnit.SECONDS)。- Semaphore 无条件
release()。凭空增加许可,限流失效。用 boolean 标记是否 acquire 成功。 - 把 Semaphore 当互斥锁用。它没有 ownership 概念,别的线程能随便
release(),而且不可重入。要互斥用ReentrantLock。 - CyclicBarrier 的线程数算错。
parties设 5 但只起了 4 个线程,第 4 个会永远等下去。我们有一次分片数量是动态算的,空列表时shards.size()是 0,new CyclicBarrier(0)直接抛IllegalArgumentException。后来加了判空。 - 在
barrierAction里做重活。它由最后一个到达的线程执行,这个线程会被拖住。对账那个任务里我在 action 里做了数据库写入,导致每批要多等 300 毫秒,4 个线程有 3 个在空转。改成只做计数统计,写库放到每轮结束后单独做。
小结
CountDownLatch是 AQS 共享模式,state 存计数,一次性,无 reset。countDown放 finally、await带超时,这两条是铁律。CyclicBarrier用 ReentrantLock + Condition,靠Generation实现可复用。任何一个线程被中断或超时,整个栅栏 broken,其他人全收到BrokenBarrierException。Semaphore没有 ownership,任何线程都能release()。所以必须判断 acquire 是否成功再 release,否则许可会越来越多。- Semaphore 公平模式吞吐只有非公平的 1/4,限流场景为了不饿死值得付这个代价。
- 我们项目里,
CountDownLatch用得最多(接口并行聚合),Semaphore次之(下游限流),CyclicBarrier只在那个对账任务里用过一次。