首页 / 资讯中心 / 文章详情

CyclicBarrier实战:从原理到多阶段线程协作的完全指南

CyclicBarrier实战:从原理到多阶段线程协作的完全指南 ★ FEATURED ARTICLE
最近有个线程协作的场景让我反复折腾了好一阵子——一批任务要分成多个阶段跑每个阶段所有线程都得齐了才能开始下一轮。一开始我想都没想就掏出了CountDownLatch结果代码越写越别扭。直到我把CyclicBarrier的源码从头到尾捋了一遍才发现之前对它的认知几乎全是听别人说的二手印象真正落地的时候那些细节根本不够用。这篇就当作一次实战复盘把CyclicBarrier的原理、代码细节和真实场景下的用法一次性说透适合那些已经会用CountDownLatch和Thread.join()、但还没真正理解CyclicBarrier独特价值的人。看完之后你应该能直接判断这个场景到底该不该用它而不是靠猜。1. 从一次失败的经历说起CountDownLatch不是万能的终局工具先说那个让我翻车的场景。当时要做的是批量数据的分阶段处理每个线程处理自己分到的数据切片处理完一轮之后必须等所有线程都处理完同一批次才能拿这批结果去做下一轮计算。第一轮、第二轮的操作类型不一样但齐头并进、统一换挡的约束是固定的。我一开始用的是CountDownLatch一个比较粗糙的做法是每个线程一轮结束就countDown()一次然后在下一轮开始前await()等计数归零。看似没问题但代码跑起来之后马上冒出两个麻烦第一每一轮我都得重新new一个CountDownLatch。因为CountDownLatch的计数是一次性的减到零它就永久失效了。为了跑五个阶段我得准备好几个实例或者用循环变量动态替换代码丑得我都不好意思看。第二假设有线程在中途因为异常被中断了它的任务没跑完但计数被减了其它线程会一直傻等永远不会等到计数归零。这算是一个比较常见的问题CountDownLatch的countDown()和await()是解耦的没人保证计数归零和所有线程真的准备好了这两件事有严格对应关系。事后复盘其实我的用法是错的我需要的不是倒计时结束之后永久放行而是每一轮都让所有参与线程在栅栏前集合一次人齐了才统一发车并且下一轮还能接着用。这正是CyclicBarrier设计的初衷——循环屏障既能等待所有线程到达同一点又能在满足条件后自动重置让下一轮继续协作。这个概念其实在生活里很常见。想想团队爬山CountDownLatch相当于在山顶立一根旗杆谁先到谁先等人齐了旗杆就退休后续谁来都不管了。而CyclicBarrier更像是一辆车每个队员到齐一个就上一个座位座位坐满后司机发车到了下一个景点大家下车休息再上车时座位重新清空。跑多阶段协作任务明显是后者的节奏。2. 原理拆解CyclicBarrier源码里到底藏着什么机关既然决定搞清楚它我直接把JDK源码打开从头读了一遍。这里把关键机制拆出来读完你会发现它其实是一个并发原语组合拳。2.1 Generation机制打破循环的钥匙CyclicBarrier里有一个内部类Generation它就做了两件事记录当前这一代栅栏是否处于损坏状态。每完成一轮完整的等待-放行流程generation就会被替换成一个新的实例。这个设计很聪明它让CyclicBarrier区分了正常情况下的一轮结束和在异常情况下的失效。如果某个线程在等待时被中断或者栅栏本身出现了BrokenBarrierExceptiongeneration会被标记为brokentrue。以后再有线程进来await()会立刻检测到broken状态然后抛出异常。这个代际设计有意思的点在于它是从根本上杜绝了上一轮的坏状态污染下一轮的可能。有一说一这个点和很多人直觉里的计数器想象差别很大。它不是简单的数字加减而是由一个状态对象把每一轮的界限划出来了。一旦某一轮坏了你不能再指望下一轮还能自动正常除非你自己决定重置或者新建实例。2.2 基于ReentrantLock的await流程排队等待的正确姿势CyclicBarrier的完整性保障来自它内部的ReentrantLock和对应的Condition。await()方法的执行逻辑大致是这样的private int dowait(boolean timed, long nanos) { final ReentrantLock lock this.lock; lock.lock(); try { final Generation g generation; if (g.broken) { throw new BrokenBarrierException(); } if (Thread.interrupted()) { breakBarrier(); throw new InterruptedException(); } int index --count; if (index 0) { // 所有线程到齐执行可选的 barrierAction然后换代 boolean ranAction false; try { final Runnable command barrierCommand; if (command ! null) { command.run(); } ranAction true; nextGeneration(); return 0; } finally { if (!ranAction) { breakBarrier(); } } } // 还没到齐进入 Condition 等待 while (true) { try { if (!timed) { trip.await(); } else if (nanos 0L) { nanos trip.awaitNanos(nanos); } } catch (InterruptedException ie) { // ... } if (g ! generation) { return index; } if (g.broken) { throw new BrokenBarrierException(); } } } finally { lock.unlock(); } }核心逻辑一句话概括把所有的参与线程通过一个锁和条件变量组织成排队等齐的协作组最后一个到达的线程触发重置和放行其余线程被逐个唤醒。这里值得注意的一个细节是Command.run()也就是barrierAction是在最后一个到达的线程里执行的。换句话说这个到了栅栏之后做点什么事情的回调并不是由什么神秘的第N1个线程来跑而是由最后一个触发放行的线程来执行的。如果你在这个回调里做了特别重的耗时操作那最后一个线程的返回会被明显拖慢但它并不会额外占用一个线程名额。2.3 对比CountDownLatch和CyclicBarrier选型其实一句话就能定维度CountDownLatchCyclicBarrier计数器来源构造时指定一个固定计数内部根据线程数维护参与计数用完后能否自动重置不能归零即失效能每轮自动换代主要等待对象等待某个或某几个线程完成一组操作等待一组线程互相到达同一汇合点是否可以重复使用不行需要重新实例化行天然支持多轮异常语义等待方可能一直阻塞有破损状态机制用异常显式通知可选回调无支持barrierAction最后一个到达线程执行只记住一条判断原则就行如果你的需求是一组线程干完各自的活之后汇聚到下一个阶段那就直接用CyclicBarrier如果你只是单纯等某个资源初始化完成或者等N个外部任务都结束那CountDownLatch更合适。3. 落地实战从一行Hello World到分批并行任务的演进原理看再多不写代码等于白看。这里我从最简单的场景一步步加深演示怎么正确使用CyclicBarrier。3.1 两个最小示例理解人齐发车和循环使用先看最典型的齐头并进用法。三个线程分别模拟三个工作人员每个阶段完成后在栅栏处等待直到三人都到达才统一进入下一阶段import java.util.concurrent.BrokenBarrierException; import java.util.concurrent.CyclicBarrier; public class BasicUsage { public static void main(String[] args) { int workerCount 3; CyclicBarrier barrier new CyclicBarrier(workerCount, () - System.out.println( 所有工人已到齐开始下一阶段 ) ); for (int i 0; i workerCount; i) { new Thread(new Worker(barrier), 工人- i).start(); } } static class Worker implements Runnable { private final CyclicBarrier barrier; Worker(CyclicBarrier barrier) { this.barrier barrier; } Override public void run() { for (int round 1; round 3; round) { // 模拟不同阶段的耗时任务 int cost 500 (round * 300); System.out.printf(%s 正在执行第 %d 阶段任务预计耗时 %dms%n, Thread.currentThread().getName(), round, cost); sleepUninterruptibly(cost); try { barrier.await(); } catch (InterruptedException | BrokenBarrierException e) { e.printStackTrace(); } } } } private static void sleepUninterruptibly(long millis) { long remaining millis; long endTime System.currentTimeMillis() millis; while (remaining 0) { try { Thread.sleep(remaining); return; } catch (InterruptedException e) { remaining endTime - System.currentTimeMillis(); } } } }这个例子看着简单但有个关键点workerCount和实际启动的线程数必须严格一致。如果你构造时传了3结果只启动了2个线程那么这2个线程会永远在await()处阻塞。反过来如果你启动了4个线程而屏障只等3个那第4个线程会直接越过屏障因为计数器在3个线程到达后就已经归零并重置了。再来看一个循环利用能带来的实际好处。比如有4个线程需要连续完成3轮任务用CyclicBarrier只需要创建一个实例public class ReuseDemo { public static void main(String[] args) { CyclicBarrier barrier new CyclicBarrier(4); for (int i 1; i 4; i) { final int id i; new Thread(() - { for (int round 1; round 3; round) { doWork(id, round); try { barrier.await(); } catch (InterruptedException | BrokenBarrierException e) { Thread.currentThread().interrupt(); return; } } }, T- id).start(); } } private static void doWork(int id, int round) { // 模拟实际工作 System.out.println(线程 id 完成第 round 轮工作); } }如果没有循环能力这段代码就得在每一轮结束时构造一个新的CountDownLatch(4)还要想办法让所有线程拿到同一个新实例——这中间就很容易出现有的线程还在等旧栅栏有的线程已经开始用新栅栏的竞态。CyclicBarrier天然规避了这类错误。3.2 分批并行任务真正的落地场景实际业务里面最典型的使用场景我总结了一下常见有三种多线程分片加载数据每片独立处理完成后汇总并行请求多个下游服务等所有请求返回后再继续迭代式算法每轮并行计算结果完全同步后再跑下一轮这几个场景有一个共性阶段之间有强同步屏障且阶段数不确定或很多。这里我以一个多线程分片统计多轮迭代的案例来示范。场景设定假设有一个大文本文件里面每行一个整数。要把所有整数做归约求和但数据量太大需要分片并行累加每轮只允许处理部分数据然后合并中间结果再进入下一轮。用CyclicBarrier的好处是每轮结束之后所有线程都必须等到其他线程的中间结果就绪才能安全地进入下一轮不会出现数据覆盖问题。import java.util.ArrayList; import java.util.List; import java.util.concurrent.*; public class SplitSumTask { private final int workerCount; private final ListInteger data; private final CyclicBarrier barrier; private volatile int globalSum 0; private final ListInteger partialSums new ArrayList(); public SplitSumTask(ListInteger data, int workerCount) { this.data data; this.workerCount workerCount; this.barrier new CyclicBarrier(workerCount, this::mergePartialResults); } public int run() throws InterruptedException, BrokenBarrierException { int partSize (data.size() workerCount - 1) / workerCount; ListThread threads new ArrayList(); for (int w 0; w workerCount; w) { final int from w * partSize; final int to Math.min((w 1) * partSize, data.size()); Thread t new Thread(() - { int localSum 0; // 只执行一轮但可以在循环里做成多轮迭代 for (int i from; i to; i) { localSum data.get(i); } synchronized (partialSums) { partialSums.add(localSum); } try { barrier.await(); } catch (InterruptedException | BrokenBarrierException e) { Thread.currentThread().interrupt(); } }, worker- w); threads.add(t); t.start(); } for (Thread t : threads) { t.join(); } return globalSum; } private void mergePartialResults() { // 这个回调在最后一个到达的线程中执行 int sum 0; synchronized (partialSums) { for (Integer v : partialSums) { sum v; } partialSums.clear(); } globalSum sum; System.out.println(阶段汇总结果 sum); } public static void main(String[] args) throws Exception { ListInteger numbers new ArrayList(); for (int i 0; i 100000; i) { numbers.add(i); } SplitSumTask task new SplitSumTask(numbers, 4); int result task.run(); System.out.println(最终总和 result); } }这段代码我故意写得直白方便看懂。几个值得注意的细节mergePartialResults里加了synchronized因为partialSums.add()这个步骤发生在await()之前而mergePartialResults执行在最后一个线程的await()内部——这两个时刻是不同的线程在执行存在可见性问题所以变量本身用volatile修饰或者放到锁保护下都需要注意。其实更优雅的写法是让每个worker自己管理一个局部变量数组合并的时候按索引读取避免额外的锁。但为了示例可读性先用同步块。如果mergePartialResults里面抛了异常会触发breakBarrier()把栅栏标记为broken。所有等待中的线程都会抛BrokenBarrierException这个事件本身是个重要的失败信号。很多人在写回调时忽略了异常处理这是个大坑。这个例子只跑了一轮工程上通常会把for循环加大让它跑多轮迭代。只要把上面代码包在一个while循环里每次循环结束之后都检查一下全局状态是否符合预期不符合就break完全能实现动态迭代次数的版本。3.3 barrierAction的三种常见用法barrierAction是我觉得CyclicBarrier被低估的一个功能。它让线程全部到达后、正式放行前这个时间窗口有了明确归属。实际用下来基本就三种玩法做汇总像上面那个例子一样最后一个线程负责把各自的中间结果合并成全局结果做发布把所有线程已完成本阶段工作这个事件通知出去比如切换数据库连接、刷新配置做校验在进入下一阶段之前检查一下上一阶段的结果是否正确不满足条件就主动破坏栅栏或做补偿有一个容易踩坑的点是barrierAction是在最后一个到达线程的调用栈上执行的。如果你的回调需要访问某些线程私有状态而这些状态还没走到合适的同步点就可能有可见性问题。建议回调里只做汇总、打印、通知这类轻量操作并且所有被访问的数据都做好同步或者干脆设计成线程安全的。4. 深入解读这五个坑我踩过你未必能避开纸上谈兵完了说几个实际项目中真正让我耗时最久的坑。这些坑说明书上不会写但只要你用CyclicBarrier做重活几乎都会碰到。4.1 中断和BrokenBarrierException的连锁反应这是最容易被忽视的。当一个线程在await()等待时被中断CyclicBarrier会认为这一轮已经破了立刻把generation标记为broken并唤醒所有正在等待的线程。这些被唤醒的线程在trip.await()内部醒来之后会发现自己等待的generation已经被标记为broken然后抛出BrokenBarrierException。所以你的catch逻辑里至少需要处理三种异常try { barrier.await(5, TimeUnit.SECONDS); } catch (InterruptedException e) { // 当前线程等待时被中断屏障可能已经broken Thread.currentThread().interrupt(); } catch (BrokenBarrierException e) { // 有其它线程被中断/超时/回调异常导致屏障破损 } catch (TimeoutException e) { // 当前线程等待超时通常会导致屏障broken }写业务代码的时候BrokenBarrierException我建议不要简单地吞掉或者只打日志而是要思考这个协作组出了异常那别的线程怎么办。因为它们全都从await()里被唤醒然后抛了同样的异常如果每个线程都做自己的重试或者退出很容易产生重复操作或状态不一致。4.2 超时之后的栅栏状态await(timeout)超时不是一个孤立事件。根据JDK源码超时之后会让该线程调用breakBarrier()把整个栅栏搞成broken状态。这带来的影响是其他所有正在等待的线程会随之抛出BrokenBarrierException即使它们自己的超时时间还没到。这个设计初看有点伤及无辜但其实是有意为之——既然你这个组已经无法按计划在预期时间内到齐与其让其他线程继续无意义地等下去不如尽早让大家都感知到失败。所以我个人的经验是不要轻易依赖超时恢复这条路一旦一个线程超时整个批次的协作基本宣告失败要做的是全面的失败处理而不是单独重试那个超时的线程。4.3 复用陷阱你以为重置了但其实还在broken状态因为CyclicBarrier支持循环使用有人会想那我failed之后reset()一把继续用不就行了。但reset()一旦调用所有正在await()等待的线程都会收到BrokenBarrierException。也就是说如果你在一个线程里默默调用reset()其他线程会一脸懵地全部失败退出。正确的姿势应该是先让所有线程感知到broken状态并退出自己的任务逻辑等到整个协作组都落在无人等待的状态下再统一调用reset()或者干脆直接new一个新的CyclicBarrier。我之前在一个定时任务的场景里踩过这个坑凌晨的批处理跑一半超时我直接在调度线程里reset()想让它重新来一轮结果其他线程瞬间全部抛异常日志刷屏且因为异常处理得不够好部分线程把已处理的半截数据写进了结果表第二天数据对账怎么都对不上。最后把方案改成检测到broken就全体退出由调度层重新创建新批次任务问题才彻底解决。4.4 参与线程数不一致得数清楚到底谁在等构造时传的数字是约定但这个约定本身不会自我验证。一个常见场景是线程池里的核心线程数固定为N但是任务被拒绝或者任务提前返回真正执行await()的线程数少于N那这N个线程就永久阻塞了。更隐蔽的情况是你用Executors.newFixedThreadPool(N)提交了N个任务但任务内部还调用了其他线程存在层层委托的情况——最后实际到达await()的线程数可能就不是N了。我给出的一个防守性写法是永远在await()外层套一个超时版本宁可超时后重试也不要永久等死// 至少保证业务线程不会无限挂起 barrier.await(30, TimeUnit.SECONDS);这一行代码解决不了根本问题但它能防止最坏的情况——线上服务直接卡死。真正根治的办法还是在任务设计时就把谁参与协作这件事想清楚并且做好参与者数量统计最好打印一行日志当前等待线程数 parties这能在出问题时节省大量排查时间。4.5 与线程池的配合只有线程池内线程在等才顺畅还有一点关于线程池的配合稍微有点反直觉。CyclicBarrier的await()是阻塞的所以它更偏好用在我们自己管理的线程或者固定大小的线程池里。如果你用的是像newCachedThreadPool这种动态大小的线程池一旦任务并发数大于池大小线程会被排队执行此时很有可能发生前N个线程执行完在await()上睡着了线程池没有额外线程执行剩下的任务这种死锁。严格说这不是CyclicBarrier的锅而是CyclicBarrier把问题暴露得更明显。实际项目中如果要和线程池深度绑定我建议用ThreadPoolExecutor 固定大小的线程数并且确保提交任务数恰好等于池大小或者保证池的大小大于等于参与协作的线程数。5. 更进一步的实践思考从会用到用得聪明你如果已经读到这里说明不仅仅是想把API背下来而是真的想在工程里用好它。那下面这几个思考值得看一下。5.1 多级屏障与分阶段流水线单个CyclicBarrier解决的是一层同步点。真实系统里可能有多层同步点比如先并行下载数据人齐了再并行解析再人齐了并行写库。这时候简单方案是创建两个CyclicBarrier实例分别控制但每个实例要仔细设计好参与线程集合避免同一组线程在两个屏障之间产生纠缠。另一个思路是把任务本身设计成有状态的状态机每个阶段到达await()前自己检查当前状态是否允许进入下一阶段。这种方式比堆多个屏障更灵活但要求对业务流程有更强的抽象能力。我个人的习惯是两个阶段的场景用两个屏障完全可以接受超过三个阶段就把任务模型改成状态机。5.2 扩展结合CompletableFuture做异步化改造CyclicBarrier是同步阻塞模型而CompletableFuture是异步回调模型。如果协作阶段里有大量IO操作全都阻塞在await()上等同于把IO等待串行化并不划算。有人做过混合方案每个工作线程内先用CompletableFuture发起异步IO然后调barrier.await()等待异步结果返回后再继续。这样做的意义在于await()等待的其实是本线程持有的异步任务也完成了而不是让线程傻傻等待IO完成。这个思路能很大程度缓解线程资源的浪费但实现的时候要注意CompletableFuture内部如果用了公共ForkJoinPool和CyclicBarrier的阻塞模型混在一起可能出现线程饥饿。所以一定要给CompletableFuture显式指定独立的线程池。5.3 自定义后置处理让失败信号更可观测前面说过BrokenBarrierException是个统称它并没有告诉你到底是谁把栅栏搞坏了。如果线上出了问题看到一堆BrokenBarrierException日志你并不能立刻定位到是哪个任务超时或中断了。所以在我自己的代码里我会在进入await()之前设置一个ThreadLocal或者记录一个当前阶段ID 线程ID 任务标识然后异常处理里把它打印出来。这样排查效率会高很多private static final ThreadLocalString TASK_MARK new ThreadLocal(); // 在任务开始前 TASK_MARK.set(batch: batchId , part: partIndex); // 捕获异常时 catch (BrokenBarrierException e) { log.error(栅栏破损, 当前任务: {}, TASK_MARK.get(), e); // 这里可以上报指标比如 MetricRegistry.counter(cb.broken).increment() }这个习惯虽然看起来多写了几行代码但线上的收益极大。尤其当你的系统里同时跑了好几组并行批处理任务时没有这个标记字段你根本分不清这堆异常是哪个业务触发的。5.4 压测数据参考最后给一个简单的压测参考。我在一个8核机器上跑过拆分求和的任务任务总量100万个整数分成8个线程并行累加阶段数100轮。用CyclicBarrier做同步的开销大概在每轮0.5ms左右也就是100轮下来同步开销总共约50ms对比数据计算本身耗时大约几百毫秒同步开销占比可以接受。但如果把线程数从8提到64单机上下文切换成本会大幅上升每一轮的等待时间明显变长。所以线程粒度不是越多越好得根据CPU核数和任务计算密度来定常规建议是CPU核数或者2倍CPU核数超过这个量级CyclicBarrier的开销就开始反噬了。6. 小复盘用了这么久我对CyclicBarrier的真实感受写到这里我对CyclicBarrier的使用心得其实已经融入上面每一段了。如果非要提炼一句话的话我会说它是一个约定感很强的并发工具前提是你要把参与者范围、超时策略和失败恢复路径都设计齐全它才能展示出那种整齐划一的协作美感。用下来最让我感觉舒服的一点是它把同步点这个概念具象化成了一个可重复使用的对象配合barrierAction和generation状态机制你能很清晰地表达业务里等大家都到这里我们再继续的意图。而最让我痛恨的一点则是它出问题时的连坐机制——一个线程超时全组遭殃。但换个角度想在强一致性的协作场景里这种连坐恰恰是最安全的。最后再补一个实操层面容易忽略的小技巧当你用CyclicBarrier时可以在构造方法里多传一个Runnable哪怕它内部只写一行日志也能让你在观测每一轮是否正常汇聚时多一个抓手。这个回调不承担任何业务逻辑纯粹作为观测点对排查多轮并发问题有奇效。如果哪天你在代码里看到一个又一个CountDownLatch被反复new出来或者看到某个await()后面跟着一套复杂的重置逻辑不妨停下来想一想——你要的可能不是那个一次性倒计时而是一道能循环使用的栅栏。
阅读完成 · 觉得有帮助?
咨询建站