本文目录导读:

这是一个完整的CyclicBarrier案例,包含基础用法和高级特性演示:
基础使用案例 - 旅行团集合
import java.util.concurrent.BrokenBarrierException;
import java.util.concurrent.CyclicBarrier;
import java.util.concurrent.TimeUnit;
public class CyclicBarrierBasicExample {
// 模拟旅行团游客
static class Tourist implements Runnable {
private String name;
private CyclicBarrier barrier;
private int arriveTime;
public Tourist(String name, CyclicBarrier barrier, int arriveTime) {
this.name = name;
this.barrier = barrier;
this.arriveTime = arriveTime;
}
@Override
public void run() {
try {
// 模拟游客到达集合点
System.out.println(name + " 出发前往集合点...");
Thread.sleep(arriveTime * 1000);
System.out.println(name + " 到达集合点,等待其他人...");
// 等待所有游客到达
barrier.await();
// 所有游客都到达后,一起出发
System.out.println(name + " 等到了所有伙伴,大家一起出发去景点!");
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
System.out.println(name + " 被中断了");
} catch (BrokenBarrierException e) {
System.out.println(name + " 发现屏障被破坏: " + e.getMessage());
}
}
}
public static void main(String[] args) {
int touristCount = 4;
CyclicBarrier barrier = new CyclicBarrier(touristCount, () -> {
System.out.println("\n=== 所有游客集合完毕,导游开始点名!===\n");
});
System.out.println("旅游团开始集合...");
// 创建游客线程
for (int i = 1; i <= touristCount; i++) {
new Thread(new Tourist("游客" + i, barrier, i * 2)).start();
}
// 注意:CyclicBarrier可以重复使用
// 模拟第二个行程:返回集合
// 这里简单演示,不再重复执行
}
}
高级案例 - 并行计算
import java.util.concurrent.BrokenBarrierException;
import java.util.concurrent.CyclicBarrier;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.LongAdder;
public class CyclicBarrierParallelSearch {
// 模拟数据分片处理
static class DataPartition implements Runnable {
private int partId;
private int[] data;
private int start;
private int end;
private CyclicBarrier barrier;
private LongAdder sum;
private AtomicInteger count;
public DataPartition(int partId, int[] data, int start, int end,
CyclicBarrier barrier, LongAdder sum, AtomicInteger count) {
this.partId = partId;
this.data = data;
this.start = start;
this.end = end;
this.barrier = barrier;
this.sum = sum;
this.count = count;
}
@Override
public void run() {
try {
// 第一阶段:计算局部和
int localSum = 0;
for (int i = start; i < end; i++) {
localSum += data[i];
if (data[i] % 2 == 0) {
count.incrementAndGet(); // 计算偶数个数
}
}
System.out.println("分片" + partId + " 计算完成,局部和: " + localSum);
// 等待所有分片完成第一阶段
barrier.await();
// 第二阶段:汇总全局结果
sum.add(localSum);
System.out.println("分片" + partId + " 已将结果汇入全局和");
} catch (InterruptedException | BrokenBarrierException e) {
e.printStackTrace();
}
}
}
public static void main(String[] args) throws InterruptedException {
// 模拟大数据集
int dataSize = 1000;
int partitions = 4;
int[] data = new int[dataSize];
// 初始化数据
for (int i = 0; i < dataSize; i++) {
data[i] = i + 1;
}
// 验证结果
int expectedSum = dataSize * (dataSize + 1) / 2; // 等差数列求和
LongAdder globalSum = new LongAdder();
AtomicInteger evenCount = new AtomicInteger();
CyclicBarrier barrier = new CyclicBarrier(partitions, () -> {
System.out.println("\n=== 所有分片完成第一阶段计算 ===");
System.out.println("第二阶段开始...\n");
});
ExecutorService executor = Executors.newFixedThreadPool(partitions);
// 创建并提交分片任务
int chunkSize = dataSize / partitions;
for (int i = 0; i < partitions; i++) {
int start = i * chunkSize;
int end = (i == partitions - 1) ? dataSize : (start + chunkSize);
executor.submit(new DataPartition(
i, data, start, end, barrier, globalSum, evenCount));
}
executor.shutdown();
executor.awaitTermination(10, TimeUnit.SECONDS);
// 输出结果
System.out.println("\n========== 计算结果 ==========");
System.out.println("数组元素总和: " + globalSum.sum());
System.out.println("预期总和: " + expectedSum);
System.out.println("偶数个数: " + evenCount.get());
System.out.println("结果验证: " + (globalSum.sum() == expectedSum ? "✓ 正确" : "✗ 错误"));
}
}
带超时和中断的案例
import java.util.concurrent.CyclicBarrier;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import java.util.concurrent.BrokenBarrierException;
public class CyclicBarrierTimeoutExample {
static class Worker implements Runnable {
private String name;
private CyclicBarrier barrier;
private int workTime;
public Worker(String name, CyclicBarrier barrier, int workTime) {
this.name = name;
this.barrier = barrier;
this.workTime = workTime;
}
@Override
public void run() {
try {
System.out.println(name + " 开始处理工作,预计 " + workTime + " 秒");
Thread.sleep(workTime * 1000);
System.out.println(name + " 完成工作,等待其他工人...");
// 等待其他线程,超时时间为5秒
barrier.await(5, TimeUnit.SECONDS);
System.out.println(name + " 所有工人完成,开始协同工作!");
} catch (TimeoutException e) {
System.out.println(name + " 等待超时! 其他工人太慢了");
System.out.println("是否破坏屏障: " + barrier.isBroken());
} catch (InterruptedException e) {
System.out.println(name + " 被中断了");
Thread.currentThread().interrupt();
} catch (BrokenBarrierException e) {
System.out.println(name + " 发现屏障被破坏: " + e.getMessage());
}
}
}
public static void main(String[] args) {
int workers = 3;
CyclicBarrier barrier = new CyclicBarrier(workers, () -> {
System.out.println("=== 所有工人都准备好了!开始协同工作 ===");
});
// 一个工人很慢,会导致超时
new Thread(new Worker("工人A", barrier, 2)).start();
new Thread(new Worker("工人B", barrier, 3)).start();
new Thread(new Worker("工人C(慢)", barrier, 8)).start(); // 超过5秒超时时间
// 注意:由于工人C太慢,工人A和B会超时
}
}
可重置的重复使用案例
import java.util.concurrent.CyclicBarrier;
import java.util.concurrent.BrokenBarrierException;
import java.util.concurrent.atomic.AtomicInteger;
public class CyclicBarrierResetExample {
static class TeamMember implements Runnable {
private String name;
private CyclicBarrier barrier;
private AtomicInteger round;
public TeamMember(String name, CyclicBarrier barrier, AtomicInteger round) {
this.name = name;
this.barrier = barrier;
this.round = round;
}
@Override
public void run() {
try {
while (round.get() <= 3) { // 进行3轮
// 模拟每轮工作
System.out.println(name + " 执行第" + round.get() + "轮工作...");
Thread.sleep(1000);
// 等待其他成员完成本轮工作
barrier.await();
if (round.get() <= 3) {
System.out.println(name + " 第" + round.get() + "轮完成,进入下一轮");
round.incrementAndGet();
}
}
System.out.println(name + " 所有轮次完成");
} catch (InterruptedException | BrokenBarrierException e) {
System.out.println(name + " 异常: " + e.getMessage());
Thread.currentThread().interrupt();
}
}
}
public static void main(String[] args) {
int members = 3;
AtomicInteger round = new AtomicInteger(1);
// 创建CyclicBarrier,可重复使用
CyclicBarrier barrier = new CyclicBarrier(members, () -> {
System.out.println("\n=== 第" + round.get() + "轮所有成员完成!===\n");
});
// 启动成员线程
for (int i = 1; i <= members; i++) {
new Thread(new TeamMember("成员" + i, barrier, round)).start();
}
// 演示CyclicBarrier可以重复使用
System.out.println("CyclicBarrier 正在处理 " + members + " 个成员的多轮任务...");
}
}
实际应用场景 - 游戏等待大厅
import java.util.concurrent.BrokenBarrierException;
import java.util.concurrent.CyclicBarrier;
import java.util.concurrent.ThreadLocalRandom;
public class GameLobbyExample {
static class Player implements Runnable {
private String playerName;
private CyclicBarrier barrier;
private int loadingTime;
public Player(String playerName, CyclicBarrier barrier) {
this.playerName = playerName;
this.barrier = barrier;
this.loadingTime = ThreadLocalRandom.current().nextInt(1, 5);
}
@Override
public void run() {
try {
// 游戏开始前加载
System.out.println(playerName + " 正在加载游戏资源...");
Thread.sleep(loadingTime * 1000);
System.out.println(playerName + " 加载完成,等待其他玩家...");
// 等待所有玩家加载完成
barrier.await();
// 游戏开始
System.out.println(playerName + " 游戏开始!进入战斗!");
} catch (InterruptedException | BrokenBarrierException e) {
System.out.println(playerName + " 连接中断");
}
}
}
public static void main(String[] args) {
int requiredPlayers = 5;
CyclicBarrier startBarrier = new CyclicBarrier(requiredPlayers, () -> {
System.out.println("\n=================================");
System.out.println("所有玩家加载完成,游戏即将开始!");
System.out.println("=================================\n");
});
System.out.println("游戏大厅等待 " + requiredPlayers + " 名玩家进入...");
System.out.println("---- 等待玩家 ----\n");
// 模拟5个玩家进入
for (int i = 1; i <= requiredPlayers; i++) {
new Thread(new Player("玩家" + i, startBarrier)).start();
}
// 注意:可以重复使用 - 等待下一场比赛
// 这里简化演示,不实际运行下一场
}
}
| 特性 | 说明 |
|---|---|
| 可重复使用 | 计数器归零后可以重新使用,适合多轮任务 |
| 等待机制 | 所有线程调用await()后自动执行barrierAction |
| 超时控制 | 支持等待超时,超时后抛出TimeoutException |
| 中断支持 | 线程可在等待时被中断 |
| 屏障损坏 | 某个线程异常会导致屏障损坏,其他线程抛出BrokenBarrierException |
| 线程池兼容 | 可以配合线程池使用,控制执行线程数 |
与CountDownLatch的区别
// CountDownLatch - 一次性使用 CountDownLatch latch = new CountDownLatch(3); latch.countDown(); // 计数减少 latch.await(); // 等待计数为0 // CyclicBarrier - 可重复使用 CyclicBarrier barrier = new CyclicBarrier(3); barrier.await(); // 等待线程到达屏障点 // barrier自动重置,可再次使用
这些案例涵盖了CyclicBarrier的常用场景,从简单的集合点等待到复杂的并行计算和游戏协同任务,实际使用时,请根据具体需求选择合适的特性。