CyclicBarrier可重复使用屏障

wen java案例 1

CyclicBarrier 可重复使用屏障

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

CyclicBarrier可重复使用屏障

核心特点:可重复使用(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());
    }
}

最佳实践

  1. 合理设置线程数:确保 CyclicBarrier 的参与线程数与实际启动的线程数一致

  2. 处理中断和异常await() 会抛出 InterruptedExceptionBrokenBarrierException

  3. 超时处理:使用 await(timeout, unit) 防止永久等待

  4. 屏障动作注意事项:屏障动作由最后一个到达的线程执行,确保该动作不会影响线程的正常执行

  5. 重置屏障:使用 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 是处理多阶段并行任务的强大工具,特别是在需要所有线程同步执行多个阶段时非常有用。

抱歉,评论功能暂时关闭!