1. 引言
在 Java 并发编程中,CountDownLatch和CyclicBarrier是两个用于线程间协调的经典工具类。它们都基于AbstractQueuedSynchronizer(AQS)实现,但设计理念和使用场景截然不同。本文将深入剖析两者的核心区别、底层源码实现,并结合实战案例展示其典型用法,帮助开发者精准选择并高效运用。
2. 核心区别概览
| 特性 | CountDownLatch | CyclicBarrier |
|---|---|---|
| 核心机制 | 一次性计数器,减到 0 时释放所有等待线程 | 可重复使用的屏障,所有线程到达屏障后一起释放 |
| 计数方向 | 递减(countDown) | 递增(await) |
| 重置能力 | 不可重置,计数为 0 后失效 | 可重置(reset)或自动重置(换代) |
| 主从角色 | 主线程等待(await)多个子线程完成任务(countDown) | 多个线程相互等待,地位对等 |
| 典型场景 | 启动前等待资源初始化、等待多个服务启动完成 | 并行计算分阶段同步、多线程数据合并 |
3. CountDownLatch 详解
3.1 核心用法
CountDownLatch通过一个计数器工作,构造时指定初始值count。线程调用countDown()使计数器减 1,调用await()的线程会阻塞直到计数器变为 0。
// 典型用法:主线程等待多个子线程完成初始化publicclassCountDownLatchDemo{publicstaticvoidmain(String[]args)throwsInterruptedException{intworkerCount=3;CountDownLatchlatch=newCountDownLatch(workerCount);for(inti=0;i<workerCount;i++){newThread(()->{try{// 模拟初始化工作Thread.sleep((long)(Math.random()*1000));System.out.println(Thread.currentThread().getName()+" 初始化完成");}catch(InterruptedExceptione){Thread.currentThread().interrupt();}finally{latch.countDown();// 任务完成,计数器减1}},"Worker-"+i).start();}System.out.println("主线程等待所有 Worker 初始化...");latch.await();// 阻塞直到计数器为0System.out.println("所有 Worker 初始化完成,主线程继续执行");}}3.2 源码解读(基于 OpenJDK 17)
CountDownLatch内部使用一个继承自AbstractQueuedSynchronizer的静态内部类Sync来实现同步。
// java.util.concurrent.CountDownLatch.SyncprivatestaticfinalclassSyncextendsAbstractQueuedSynchronizer{Sync(intcount){setState(count);// 使用 AQS 的 state 存储计数器}intgetCount(){returngetState();}// 尝试获取共享锁,只有当 state == 0 时才成功(即计数器为0)protectedinttryAcquireShared(intacquires){return(getState()==0)?1:-1;}// 尝试释放共享锁,即执行 countDown,将 state 减1protectedbooleantryReleaseShared(intreleases){// 自旋 CAS 减1for(;;){intc=getState();if(c==0)returnfalse;// 已经为0,无法再减intnextc=c-1;if(compareAndSetState(c,nextc))returnnextc==0;// 返回 true 表示本次减1后计数器变为0}}}关键点:
- 状态(state)即计数器:构造时
setState(count)。 - await() 原理:调用
await()会触发acquireSharedInterruptibly(1),最终调用tryAcquireShared。只要state != 0,就返回 -1,导致线程进入 AQS 队列等待。 - countDown() 原理:调用
countDown()会触发releaseShared(1),最终调用tryReleaseShared。通过 CAS 将 state 减 1,如果减后变为 0,则返回 true,这会唤醒所有在await()上等待的线程。
3.3 实战场景
- 服务启动等待:等待所有微服务健康检查通过后,再启动网关。
- 并行任务汇总:MapReduce 模型中,等待所有 Map 任务完成后再启动 Reduce 任务。
- 测试并发:模拟高并发场景,让所有线程同时开始执行。
4. CyclicBarrier 详解
4.1 核心用法
CyclicBarrier允许一组线程相互等待,直到所有线程都到达某个公共屏障点(barrier)后,再一起继续执行。构造时可指定参与线程数parties以及可选的屏障动作barrierAction。
// 典型用法:多线程分阶段计算,等待所有线程完成当前阶段publicclassCyclicBarrierDemo{publicstaticvoidmain(String[]args){intthreadCount=3;CyclicBarrierbarrier=newCyclicBarrier(threadCount,()->{System.out.println("所有线程已到达屏障,开始下一阶段");});for(inti=0;i<threadCount;i++){newThread(()->{try{System.out.println(Thread.currentThread().getName()+" 开始第一阶段");Thread.sleep((long)(Math.random()*1000));barrier.await();// 等待其他线程System.out.println(Thread.currentThread().getName()+" 开始第二阶段");Thread.sleep((long)(Math.random()*1000));barrier.await();// 再次等待System.out.println(Thread.currentThread().getName()+" 完成所有阶段");}catch(InterruptedException|BrokenBarrierExceptione){Thread.currentThread().interrupt();}},"Thread-"+i).start();}}}4.2 源码解读(基于 OpenJDK 17)
CyclicBarrier内部使用ReentrantLock和Condition实现同步,核心状态包括:
parties:屏障的参与线程数count:当前尚未到达屏障的线程数(递减)generation:代表当前“代”,每次屏障被打破或重置时创建新 generation
// java.util.concurrent.CyclicBarrier 核心方法 await()publicintawait()throwsInterruptedException,BrokenBarrierException{try{returndowait(false,0L);}catch(TimeoutExceptiontoe){thrownewError(toe);// cannot happen}}privateintdowait(booleantimed,longnanos)throwsInterruptedException,BrokenBarrierException,TimeoutException{finalReentrantLocklock=this.lock;lock.lock();try{finalGenerationg=generation;if(g.broken)thrownewBrokenBarrierException();if(Thread.interrupted()){breakBarrier();thrownewInterruptedException();}intindex=--count;// 当前线程到达,计数器减1if(index==0){// 最后一个线程到达booleanranAction=false;try{finalRunnablecommand=barrierCommand;if(command!=null)command.run();// 执行屏障动作ranAction=true;nextGeneration();// 重置屏障,唤醒所有等待线程return0;}finally{if(!ranAction)breakBarrier();}}// 不是最后一个线程,进入等待for(;;){try{if(!timed)trip.await();// 在 Condition 上等待elseif(nanos>0L)nanos=trip.awaitNanos(nanos);}catch(InterruptedExceptionie){if(g==generation&&!g.broken){breakBarrier();throwie;}else{Thread.currentThread().interrupt();}}if(g.broken)thrownewBrokenBarrierException();if(g!=generation)// 屏障已换代,返回到达序号returnindex;}}finally{lock.unlock();}}// 创建新一代屏障,唤醒所有等待线程privatevoidnextGeneration(){trip.signalAll();// 唤醒所有在 Condition 上等待的线程count=parties;// 重置计数器generation=newGeneration();// 创建新代}关键点:
- 可重用性:当所有线程到达屏障后,
nextGeneration()会重置count = parties并创建新的generation,屏障可重复使用。 - 屏障动作:最后一个到达的线程执行
barrierCommand(如果存在),然后才唤醒其他线程。 - 中断处理:线程在等待期间被中断会调用
breakBarrier(),将当前 generation 标记为 broken,并唤醒所有等待线程。 - 超时机制:
await(long timeout, TimeUnit unit)支持超时,超时后屏障被打破。
4.3 实战场景
- 并行计算分阶段同步:如多线程排序算法,每完成一个阶段后同步数据。
- 多线程数据合并:多个线程分别处理数据的一部分,全部完成后合并结果。
- 游戏服务器同步:多个玩家准备就绪后同时开始游戏。
- 批量任务处理:将大任务拆分为多个子任务,所有子任务完成后执行汇总操作。
5. 对比总结与选型建议
5.1 核心差异对比
| 维度 | CountDownLatch | CyclicBarrier |
|---|---|---|
| 重用性 | 一次性,计数为0后失效 | 可重复使用,自动或手动重置 |
| 计数方向 | 递减(countDown) | 递减(内部count),但逻辑上是递增到达 |
| 线程角色 | 主从模式(1个主线程等待N个子线程) | 对等模式(N个线程相互等待) |
| 屏障动作 | 无 | 支持(最后一个线程到达后执行) |
| 异常处理 | 计数不会重置 | 中断或超时会导致屏障被打破(BrokenBarrierException) |
| 适用场景 | 一等多、任务完成后触发 | 多等多、分阶段同步 |
5.2 选型指南
选择 CountDownLatch 当:
- 需要一次性事件通知机制
- 主线程需要等待多个子线程完成初始化或准备工作
- 不需要重复使用同步点
- 例如:服务启动等待、测试并发开始信号
选择 CyclicBarrier 当:
- 需要多个线程在某个点同步后继续执行
- 需要重复使用同步屏障
- 需要在所有线程到达后执行特定操作(屏障动作)
- 例如:并行计算分阶段、多轮游戏同步、批量数据处理
5.3 混合使用示例
在实际项目中,两者可以结合使用:
// 使用 CountDownLatch 等待所有线程初始化完成// 使用 CyclicBarrier 进行多轮计算同步publicclassHybridDemo{publicstaticvoidmain(String[]args)throwsInterruptedException{intthreadCount=4;intphases=3;CountDownLatchinitLatch=newCountDownLatch(threadCount);CyclicBarrierphaseBarrier=newCyclicBarrier(threadCount);for(inti=0;i<threadCount;i++){newThread(()->{// 初始化阶段System.out.println(Thread.currentThread().getName()+" 初始化完成");initLatch.countDown();try{initLatch.await();// 等待所有线程初始化完成// 多阶段计算for(intphase=1;phase<=phases;phase++){System.out.println(Thread.currentThread().getName()+" 开始第 "+phase+" 阶段");Thread.sleep((long)(Math.random()*500));phaseBarrier.await();// 等待其他线程完成当前阶段}System.out.println(Thread.currentThread().getName()+" 所有阶段完成");}catch(Exceptione){Thread.currentThread().interrupt();}},"Worker-"+i).start();}}}6. 常见问题与注意事项
6.1 CountDownLatch 常见问题
- 计数溢出:
countDown()调用次数超过初始计数不会报错,但可能导致逻辑错误。 - 不可重置:计数为0后无法重用,需要创建新实例。
- 线程安全:
countDown()和await()本身是线程安全的,但业务逻辑需要自行保证。 - 超时等待:使用
await(long timeout, TimeUnit unit)避免永久阻塞。
6.2 CyclicBarrier 常见问题
- 屏障被打破:线程中断、超时或屏障动作异常会导致屏障被打破,所有等待线程收到
BrokenBarrierException。 - 重置风险:调用
reset()会打破当前屏障,可能导致等待线程收到异常。 - 屏障动作异常:屏障动作抛出异常会导致屏障被打破。
- 线程数匹配:构造时指定的
parties必须与实际调用await()的线程数一致。
6.3 性能考量
- CountDownLatch:基于 AQS,适合一次性同步场景,性能开销较小。
- CyclicBarrier:基于
ReentrantLock和Condition,适合可重复使用的同步场景,但每次屏障换代都有一定开销。 - 在高并发场景下,考虑使用
Phaser(Java 7+)作为更灵活的替代方案。
7. 总结
CountDownLatch和CyclicBarrier都是 Java 并发包中强大的线程协调工具,但设计理念不同:
CountDownLatch是一次性的倒数门闩,适合一等多的场景,如主线程等待多个子线程完成任务。CyclicBarrier是可循环使用的屏障,适合多等多的场景,如多个线程需要分阶段同步。
选择时需根据业务场景的同步需求、重用性要求和线程角色来决定。理解它们的源码实现有助于避免常见的并发陷阱,编写出更健壮、高效的并发程序。
在实际开发中,还可以结合CompletableFuture、Phaser等更现代的并发工具,构建更复杂的同步模式。