CyclicBarrier 可重复使用屏障
CyclicBarrier 是 Java 并发包 (java.util.concurrent) 中的一个同步辅助类,它允许一组线程互相等待,直到所有线程都到达一个共同的屏障点(barrier point)后再继续执行。

核心特点:可重复使用(Cyclic)
"Cyclic"(循环的)意味着当所有线程都越过一次屏障后,屏障会自动重置,可以再次使用,这是它与 CountDownLatch 最大的区别。
基本工作原理
public class CyclicBarrierDemo {
public static void main(String[] args) {
// 创建一个需要3个线程到达的屏障
CyclicBarrier barrier = new CyclicBarrier(3, () -> {
System.out.println("所有线程到达屏障,执行屏障动作!");
});
for (int i = 0; i < 3; i++) {
new Thread(() -> {
try {
System.out.println(Thread.currentThread().getName() + " 开始工作");
Thread.sleep(1000);
System.out.println(Thread.currentThread().getName() + " 到达屏障");
barrier.await(); // 等待其他线程
System.out.println(Thread.currentThread().getName() + " 越过屏障继续执行");
} catch (Exception e) {
e.printStackTrace();
}
}, "Thread-" + i).start();
}
}
}
输出示例:
Thread-0 开始工作
Thread-1 开始工作
Thread-2 开始工作
Thread-0 到达屏障
Thread-1 到达屏障
Thread-2 到达屏障
所有线程到达屏障,执行屏障动作!
Thread-2 越过屏障继续执行
Thread-1 越过屏障继续执行
Thread-0 越过屏障继续执行
关键特性
可重复使用
public class CyclicBarrierReuseDemo {
public static void main(String[] args) {
CyclicBarrier barrier = new CyclicBarrier(3);
for (int round = 0; round < 3; round++) { // 同一屏障使用3次
System.out.println("\n=== 第 " + (round + 1) + " 轮 ===");
for (int i = 0; i < 3; i++) {
new Thread(() -> {
try {
Thread.sleep((long)(Math.random() * 1000));
System.out.println(Thread.currentThread().getName() + " 到达屏障");
barrier.await(); // 每次都使用同一个 barrier
System.out.println(Thread.currentThread().getName() + " 继续执行");
} catch (Exception e) {
e.printStackTrace();
}
}, "Thread-" + i).start();
}
// 等待所有线程完成一轮
try {
Thread.sleep(2000);
} catch (InterruptedException e) {
e.printStackTrace();
}
}
}
}
屏障动作(Barrier Action)
当最后一个线程到达时,可以执行一个 Runnable,这个动作由最后一个到达的线程执行。
CyclicBarrier barrier = new CyclicBarrier(4, () -> {
System.out.println("所有选手就位,比赛开始!");
});
等待超时
public class TimeoutDemo {
public static void main(String[] args) {
CyclicBarrier barrier = new CyclicBarrier(3);
new Thread(() -> {
try {
// 等待2秒,超时则抛出 TimeoutException
barrier.await(2, TimeUnit.SECONDS);
} catch (TimeoutException e) {
System.out.println("等待超时!");
} catch (Exception e) {
e.printStackTrace();
}
}).start();
}
}
与 CountDownLatch 的对比
| 特性 | CyclicBarrier | CountDownLatch |
|---|---|---|
| 重用性 | ✅ 可重复使用 | ❌ 一次性使用 |
| 触发机制 | 所有线程到达后释放 | 计数器归零后释放 |
| 参与者 | 线程互相等待 | 线程等待计数归零 |
| 屏障动作 | 支持(可选 Runnable) | 不支持 |
| 中断 | 支持 | 支持 |
| 超时 | 支持 | 支持 |
实际应用场景
并行计算 - 分阶段汇总
public class ParallelComputation {
private static final int THREAD_COUNT = 4;
public static void main(String[] args) {
CyclicBarrier barrier = new CyclicBarrier(THREAD_COUNT, () -> {
System.out.println("所有线程完成当前阶段,进入下一阶段");
});
ExecutorService executor = Executors.newFixedThreadPool(THREAD_COUNT);
for (int i = 0; i < THREAD_COUNT; i++) {
executor.submit(() -> {
try {
// 阶段1:数据加载
System.out.println(Thread.currentThread().getName() + " 加载数据完成");
barrier.await();
// 阶段2:数据处理
System.out.println(Thread.currentThread().getName() + " 处理数据完成");
barrier.await();
// 阶段3:结果输出
System.out.println(Thread.currentThread().getName() + " 输出结果完成");
} catch (Exception e) {
e.printStackTrace();
}
});
}
executor.shutdown();
}
}
游戏开发 - 多玩家加载
public class GameLobby {
private CyclicBarrier barrier;
private int playerCount;
public GameLobby(int playerCount) {
this.playerCount = playerCount;
this.barrier = new CyclicBarrier(playerCount, () -> {
System.out.println("所有玩家加载完成,游戏开始!");
});
}
public void loadGame(String playerName) {
new Thread(() -> {
try {
System.out.println(playerName + " 正在加载游戏...");
Thread.sleep(1000); // 模拟加载
System.out.println(playerName + " 加载完成,等待其他玩家");
barrier.await();
System.out.println(playerName + " 进入游戏!");
} catch (Exception e) {
e.printStackTrace();
}
}).start();
}
}
故障处理机制
public class BrokenBarrierDemo {
public static void main(String[] args) {
CyclicBarrier barrier = new CyclicBarrier(3);
// 线程1正常等待
new Thread(() -> {
try {
System.out.println("线程1 等待");
barrier.await();
} catch (Exception e) {
System.out.println("线程1 中断: " + e.getMessage());
}
}).start();
// 线程2中断等待
Thread t2 = new Thread(() -> {
try {
System.out.println("线程2 等待");
Thread.sleep(100);
barrier.await();
} catch (Exception e) {
System.out.println("线程2 异常: " + e.getMessage());
}
});
t2.start();
// 中断线程2
t2.interrupt();
// 检查屏障状态
System.out.println("屏障是否损坏: " + barrier.isBroken());
}
}
最佳实践
-
合理设置线程数:确保
CyclicBarrier的参与线程数与实际启动的线程数一致 -
处理中断和异常:
await()会抛出InterruptedException和BrokenBarrierException -
超时处理:使用
await(timeout, unit)防止永久等待 -
屏障动作注意事项:屏障动作由最后一个到达的线程执行,确保该动作不会影响线程的正常执行
-
重置屏障:使用
reset()方法可以手动重置屏障(会抛出BrokenBarrierException)
// 安全使用示例
public void safeUse(CyclicBarrier barrier) {
try {
barrier.await(5, TimeUnit.SECONDS);
} catch (TimeoutException e) {
// 处理超时
barrier.reset(); // 重置屏障
} catch (BrokenBarrierException e) {
// 屏障被破坏
} catch (InterruptedException e) {
Thread.currentThread().interrupt(); // 恢复中断状态
}
}
CyclicBarrier 是处理多阶段并行任务的强大工具,特别是在需要所有线程同步执行多个阶段时非常有用。