Java Phaser案例

wen java案例 2

本文目录导读:

Java Phaser案例

  1. 基础案例:多阶段任务协同
  2. 高级案例:模拟运动会比赛
  3. 实用案例:流水线数据处理系统
  4. 动态注册/注销案例
  5. 中断和异常处理案例

我将为您提供一个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的核心功能:

  1. 多阶段同步:参与者可以在多个阶段进行同步
  2. 动态注册:可以动态添加或删除参与者
  3. 复用性:Phaser可以被重复使用
  4. 灵活性:支持中断、超时等操作
  5. 批量操作:支持arriveAndAwaitAdvance等原子操作

在实际应用中,Phaser非常适合处理需要多阶段协作的场景,如模拟测试、并发数据流水线处理、多阶段算法实现等。

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