本文目录导读:

我将为您提供一个Java Phaser的完整案例,展示其强大的多阶段同步功能。
基础案例:多阶段任务协同
import java.util.concurrent.Phaser;
public class PhaserBasicExample {
public static void main(String[] args) {
// 创建Phaser,初始parties = 1(主线程)
Phaser phaser = new Phaser(1);
// 创建3个工作线程
for (int i = 0; i < 3; i++) {
Worker worker = new Worker("Worker-" + i, phaser);
Thread thread = new Thread(worker);
thread.start();
}
// 主线程执行3个阶段
for (int phase = 0; phase < 3; phase++) {
System.out.println("=== 主线程启动阶段 " + phase + " ===");
// 主线程到达并等待其他线程完成当前阶段
phaser.arriveAndAwaitAdvance();
System.out.println("=== 阶段 " + phase + " 完成,当前phase="
+ phaser.getPhase() + " ===");
System.out.println();
}
// 注销主线程
phaser.arriveAndDeregister();
System.out.println("所有阶段完成!");
System.out.println("Phaser终止状态: " + phaser.isTerminated());
}
static class Worker implements Runnable {
private final String name;
private final Phaser phaser;
public Worker(String name, Phaser phaser) {
this.name = name;
this.phaser = phaser;
// 注册工作线程
phaser.register();
}
@Override
public void run() {
// 执行3个阶段的工作
for (int phase = 0; phase < 3; phase++) {
System.out.println(name + " 开始阶段 " + phase + " 的工作");
// 模拟工作
try {
Thread.sleep((long) (Math.random() * 1000));
} catch (InterruptedException e) {
e.printStackTrace();
}
System.out.println(name + " 完成阶段 " + phase + " 的工作");
// 到达并等待其他线程
phaser.arriveAndAwaitAdvance();
}
// 注销线程
phaser.arriveAndDeregister();
System.out.println(name + " 完成所有工作并注销");
}
}
}
高级案例:模拟运动会比赛
import java.util.concurrent.Phaser;
import java.util.concurrent.ThreadLocalRandom;
public class SportsGameExample {
private static final int PLAYERS = 4;
private static final int ROUNDS = 3;
public static void main(String[] args) {
System.out.println("===== 运动会比赛开始 =====");
// 创建Phaser,初始parties = 1(裁判)
Phaser referee = new Phaser(1);
// 创建运动员
for (int i = 1; i <= PLAYERS; i++) {
new Thread(new Athlete("运动员" + i, referee)).start();
}
// 裁判主持比赛
for (int round = 1; round <= ROUNDS; round++) {
System.out.println("\n--- 第 " + round + " 轮比赛开始 ---");
// 裁判宣布比赛开始
referee.arriveAndAwaitAdvance();
// 等待所有运动员完成比赛
referee.arriveAndAwaitAdvance();
System.out.println("--- 第 " + round + " 轮比赛结束 ---");
}
// 结束比赛
referee.arriveAndDeregister();
System.out.println("\n===== 比赛结束 =====");
}
static class Athlete implements Runnable {
private final String name;
private final Phaser phaser;
public Athlete(String name, Phaser phaser) {
this.name = name;
this.phaser = phaser;
phaser.register();
}
@Override
public void run() {
for (int round = 1; round <= ROUNDS; round++) {
// 等待比赛开始
phaser.arriveAndAwaitAdvance();
// 进行比赛
int time = ThreadLocalRandom.current().nextInt(500, 2000);
System.out.println(name + " 正在比赛,预计用时: " + time + "ms");
try {
Thread.sleep(time);
} catch (InterruptedException e) {
e.printStackTrace();
}
System.out.println(name + " 完成第 " + round + " 轮比赛");
// 报告完成
phaser.arriveAndAwaitAdvance();
}
phaser.arriveAndDeregister();
System.out.println(name + " 完成所有比赛");
}
}
}
实用案例:流水线数据处理系统
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;
public class PipelineExample {
private static final int STAGES = 4;
private static final int DATA_COUNT = 10;
private static final AtomicInteger stage2Result = new AtomicInteger(0);
public static void main(String[] args) {
System.out.println("===== 流水线数据处理系统 =====");
// 创建4阶段Phaser
Phaser phaser = new Phaser(1);
// 各个阶段的数据处理器
StageProcessor[] processors = new StageProcessor[STAGES];
Thread[] threads = new Thread[STAGES];
for (int i = 0; i < STAGES; i++) {
processors[i] = new StageProcessor("Stage-" + (i + 1), phaser, i);
threads[i] = new Thread(processors[i]);
threads[i].start();
}
// 主线程模拟数据产生
for (int i = 0; i < DATA_COUNT; i++) {
System.out.println("\n>>> 主线程提交数据 " + i);
phaser.arriveAndAwaitAdvance(); // 每个阶段的数据传输
// 模拟数据流经整个流水线
for (int phase = 0; phase < STAGES - 1; phase++) {
phaser.arriveAndAwaitAdvance();
}
}
// 关闭流水线
for (int i = 0; i < STAGES; i++) {
phaser.arriveAndDeregister();
}
System.out.println("\n===== 流水线处理完成 =====");
System.out.println("总共处理数据量: " + stage2Result.get());
}
static class StageProcessor implements Runnable {
private final String name;
private final Phaser phaser;
private final int stageNumber;
private int processedCount = 0;
public StageProcessor(String name, Phaser phaser, int stageNumber) {
this.name = name;
this.phaser = phaser;
this.stageNumber = stageNumber;
phaser.register();
}
@Override
public void run() {
while (!phaser.isTerminated()) {
// 等待数据到达
phaser.arriveAndAwaitAdvance();
// 处理数据
if (!phaser.isTerminated()) {
processData();
}
}
System.out.println(name + " 停止处理");
}
private void processData() {
processedCount++;
// 模拟不同的处理时间
int delay = (stageNumber + 1) * 200;
try {
Thread.sleep(delay);
} catch (InterruptedException e) {
e.printStackTrace();
}
System.out.println(name + " 处理数据 #" + processedCount);
if (stageNumber == 1) { // 第二阶段累计计数
stage2Result.incrementAndGet();
}
}
}
}
动态注册/注销案例
import java.util.concurrent.Phaser;
import java.util.concurrent.atomic.AtomicInteger;
public class DynamicRegistrationExample {
private static final int INITIAL_PARTIES = 3;
private static final int MAX_PARTIES = 6;
public static void main(String[] args) {
System.out.println("===== 动态注册Phaser示例 =====");
Phaser phaser = new Phaser(INITIAL_PARTIES);
// 初始3个参与者
for (int i = 0; i < INITIAL_PARTIES; i++) {
new Thread(new Task("任务-" + (i + 1), phaser)).start();
}
// 模拟动态添加参与者
try {
Thread.sleep(1000);
System.out.println("\n添加新的参与者...");
for (int i = INITIAL_PARTIES; i < MAX_PARTIES; i++) {
new Thread(new Task("动态任务-" + (i + 1), phaser)).start();
Thread.sleep(500);
}
} catch (InterruptedException e) {
e.printStackTrace();
}
// 等待所有任务完成
try {
Thread.sleep(5000);
} catch (InterruptedException e) {
e.printStackTrace();
}
System.out.println("\n当前registered parties: " + phaser.getRegisteredParties());
System.out.println("当前phase: " + phaser.getPhase());
System.out.println("===== 示例结束 =====");
}
static class Task implements Runnable {
private final String name;
private final Phaser phaser;
public Task(String name, Phaser phaser) {
this.name = name;
this.phaser = phaser;
phaser.register();
System.out.println(name + " 注册,当前parties: " + phaser.getRegisteredParties());
}
@Override
public void run() {
try {
for (int phase = 0; phase < 3; phase++) {
System.out.println(name + " 执行阶段 " + phase);
Thread.sleep(1000);
// 当前阶段完成后注销
if (phase == 1 && name.startsWith("动态")) {
System.out.println(name + " 在阶段 " + phase + " 后注销");
phaser.arriveAndDeregister();
System.out.println(name + " 已注销,当前parties: " + phaser.getRegisteredParties());
return;
}
phaser.arriveAndAwaitAdvance();
}
// 完成所有阶段后注销
phaser.arriveAndDeregister();
} catch (InterruptedException e) {
e.printStackTrace();
}
}
}
}
中断和异常处理案例
import java.util.concurrent.Phaser;
import java.util.concurrent.TimeUnit;
public class ExceptionHandlingExample {
public static void main(String[] args) {
System.out.println("===== 中断和异常处理示例 =====");
Phaser phaser = new Phaser(3);
// 正常任务
Thread normalTask = new Thread(() -> {
System.out.println("正常任务开始");
try {
Thread.sleep(2000);
phaser.arriveAndAwaitAdvance();
System.out.println("正常任务完成");
} catch (InterruptedException e) {
e.printStackTrace();
}
});
// 会抛异常的任务
Thread errorTask = new Thread(() -> {
System.out.println("异常任务开始");
try {
Thread.sleep(1000);
throw new RuntimeException("任务出错!");
} catch (InterruptedException e) {
e.printStackTrace();
} finally {
// 确保在出错时也能正确注销
phaser.arriveAndDeregister();
System.out.println("异常任务注销");
}
});
Thread normalTask2 = new Thread(() -> {
System.out.println("第二个正常任务开始");
try {
Thread.sleep(3000);
phaser.arriveAndAwaitAdvance();
System.out.println("第二个正常任务完成");
} catch (InterruptedException e) {
e.printStackTrace();
}
});
normalTask.start();
errorTask.start();
normalTask2.start();
// 等待所有线程完成
try {
normalTask.join();
errorTask.join();
normalTask2.join();
} catch (InterruptedException e) {
e.printStackTrace();
}
System.out.println("当前parties: " + phaser.getRegisteredParties());
System.out.println("===== 结束 =====");
}
}
这些案例展示了Phaser的核心功能:
- 多阶段同步:参与者可以在多个阶段进行同步
- 动态注册:可以动态添加或删除参与者
- 复用性:Phaser可以被重复使用
- 灵活性:支持中断、超时等操作
- 批量操作:支持arriveAndAwaitAdvance等原子操作
在实际应用中,Phaser非常适合处理需要多阶段协作的场景,如模拟测试、并发数据流水线处理、多阶段算法实现等。