Java CountDownLatch案例

wen java案例 3

本文目录导读:

Java CountDownLatch案例

  1. 基础案例:主线程等待多个子线程完成
  2. 模拟并发启动:所有线程同时开始
  3. 生产环境实战:多线程数据汇总
  4. 带超时控制的案例
  5. 使用规范和最佳实践

我来为你提供几个完整的Java CountDownLatch案例,从基础到进阶。

基础案例:主线程等待多个子线程完成

import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
public class BasicCountDownLatchDemo {
    public static void main(String[] args) throws InterruptedException {
        // 创建一个计数器为3的CountDownLatch
        CountDownLatch latch = new CountDownLatch(3);
        System.out.println("主线程开始执行...");
        // 启动3个子线程
        for (int i = 1; i <= 3; i++) {
            final int taskId = i;
            new Thread(() -> {
                try {
                    System.out.println("任务" + taskId + " 开始执行");
                    TimeUnit.SECONDS.sleep(taskId); // 模拟不同耗时
                    System.out.println("任务" + taskId + " 执行完成");
                } catch (InterruptedException e) {
                    e.printStackTrace();
                } finally {
                    latch.countDown(); // 计数器减1
                }
            }).start();
        }
        System.out.println("主线程等待所有子任务完成...");
        latch.await(); // 等待计数器归零
        System.out.println("所有任务完成,主线程继续执行");
    }
}

模拟并发启动:所有线程同时开始

import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
public class ConcurrentStartDemo {
    public static void main(String[] args) throws InterruptedException {
        // 使用两个CountDownLatch实现并发启动
        CountDownLatch readyLatch = new CountDownLatch(5); // 等待5个线程就绪
        CountDownLatch startLatch = new CountDownLatch(1); // 控制启动信号
        // 创建5个工作线程
        for (int i = 1; i <= 5; i++) {
            final int workerId = i;
            new Thread(() -> {
                try {
                    // 模拟准备工作
                    System.out.println("工作者" + workerId + " 正在准备...");
                    TimeUnit.SECONDS.sleep(1);
                    System.out.println("工作者" + workerId + " 准备就绪");
                    readyLatch.countDown(); // 告知主线程已就绪
                    // 等待启动信号
                    startLatch.await();
                    System.out.println("工作者" + workerId + " 开始工作");
                    // 模拟工作
                    TimeUnit.SECONDS.sleep(2);
                    System.out.println("工作者" + workerId + " 工作完成");
                } catch (InterruptedException e) {
                    e.printStackTrace();
                }
            }).start();
        }
        // 等待所有工作者就绪
        System.out.println("主线程等待所有工作者就绪...");
        readyLatch.await();
        System.out.println("所有工作者已就绪,准备同时启动");
        // 释放启动信号,让所有工作者同时开始
        Thread.sleep(1000);
        startLatch.countDown();
        System.out.println("启动信号已发出");
    }
}

生产环境实战:多线程数据汇总

import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.*;
import java.util.Random;
public class DataAggregationDemo {
    static class DataResult {
        private final int workerId;
        private final int dataCount;
        private final long totalData;
        public DataResult(int workerId, int dataCount, long totalData) {
            this.workerId = workerId;
            this.dataCount = dataCount;
            this.totalData = totalData;
        }
        @Override
        public String toString() {
            return String.format("工作者%d:处理了%d条数据,数据总量为%d", 
                workerId, dataCount, totalData);
        }
    }
    public static void main(String[] args) throws InterruptedException {
        final int workerCount = 5;
        CountDownLatch doneLatch = new CountDownLatch(workerCount);
        List<DataResult> results = new CopyOnWriteArrayList<>();
        ExecutorService executor = Executors.newFixedThreadPool(workerCount);
        Random random = new Random();
        System.out.println("开始并行处理数据...");
        // 启动多个工作线程
        for (int i = 1; i <= workerCount; i++) {
            final int workerId = i;
            executor.submit(() -> {
                try {
                    // 模拟从数据库或其他系统读取数据
                    int dataCount = random.nextInt(100) + 50;
                    long totalData = 0;
                    // 模拟数据处理
                    for (int j = 0; j < dataCount; j++) {
                        totalData += random.nextInt(1000);
                        Thread.sleep(10); // 模拟耗时
                    }
                    results.add(new DataResult(workerId, dataCount, totalData));
                    System.out.println("工作者" + workerId + " 处理完成");
                } catch (InterruptedException e) {
                    e.printStackTrace();
                } finally {
                    doneLatch.countDown();
                }
            });
        }
        // 主线程等待所有任务完成
        doneLatch.await(10, TimeUnit.SECONDS); // 最多等待10秒
        System.out.println("\n===== 数据处理汇总 =====");
        results.forEach(System.out::println);
        // 计算汇总
        int totalCount = results.stream().mapToInt(r -> r.dataCount).sum();
        long totalData = results.stream().mapToLong(r -> r.totalData).sum();
        System.out.println("\n总计处理数据条数:" + totalCount);
        System.out.println("总计数据量:" + totalData);
        executor.shutdown();
    }
}

带超时控制的案例

import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
public class TimeoutDemo {
    public static void main(String[] args) {
        CountDownLatch latch = new CountDownLatch(2);
        // 第一个任务会成功完成
        new Thread(() -> {
            try {
                System.out.println("任务1开始执行");
                TimeUnit.SECONDS.sleep(2);
                System.out.println("任务1完成");
                latch.countDown();
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
        }).start();
        // 第二个任务模拟超时
        new Thread(() -> {
            try {
                System.out.println("任务2开始执行");
                TimeUnit.SECONDS.sleep(10); // 耗时的任务
                System.out.println("任务2完成");
                latch.countDown();
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
        }).start();
        try {
            // 等待最多3秒,如果任务未完成则超时继续
            boolean completed = latch.await(3, TimeUnit.SECONDS);
            if (completed) {
                System.out.println("所有任务在3秒内完成");
            } else {
                System.out.println("等待超时!部分任务可能还没完成");
            }
        } catch (InterruptedException e) {
            e.printStackTrace();
        }
    }
}

使用规范和最佳实践

import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
public class BestPracticeDemo {
    public static void main(String[] args) throws InterruptedException {
        ExecutorService executor = Executors.newFixedThreadPool(4);
        CountDownLatch latch = new CountDownLatch(4);
        for (int i = 0; i < 4; i++) {
            final int taskId = i;
            executor.submit(() -> {
                try {
                    try {
                        // 业务逻辑
                        System.out.println("任务" + taskId + "开始");
                        Thread.sleep(1000);
                        System.out.println("任务" + taskId + "结束");
                    } catch (InterruptedException e) {
                        Thread.currentThread().interrupt(); // 恢复中断状态
                    } finally {
                        // 确保在 finally 中调用 countDown
                        latch.countDown();
                    }
                } catch (Exception e) {
                    // 处理业务异常
                    System.err.println("任务" + taskId + "异常:" + e.getMessage());
                }
            });
        }
        // 等待完成
        boolean completed = latch.await(10, TimeUnit.SECONDS);
        if (completed) {
            System.out.println("所有任务完成");
        } else {
            System.out.println("任务超时");
        }
        executor.shutdown();
    }
}
  1. countDown() 必须放在 finally 中:确保即使出现异常,计数器也能正确递减,避免死锁。

  2. 使用带超时的 await()latch.await(timeout, timeUnit) 可以防止永久等待。

  3. 初始化计数器的值:应该等于需要等待完成的线程或操作数量。

  4. 线程安全:CountDownLatch 是线程安全的,可以被多个线程安全地调用。

  5. 不可重用:CountDownLatch 是一次性工具,计数归零后不能再使用,如果需要循环使用,应该考虑使用 CyclicBarrier。

  6. 应用场景

    • 主线程等待多个子线程完成
    • 模拟高并发场景
    • 任务拆分汇总
    • 服务启动时等待依赖组件就绪

这就是 CountDownLatch 的常用案例,涵盖了从基础到实战的多个场景,你可以根据自己的需求选择合适的案例进行参考。

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