本文目录导读:

我来为你提供几个完整的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();
}
}
-
countDown() 必须放在 finally 中:确保即使出现异常,计数器也能正确递减,避免死锁。
-
使用带超时的 await():
latch.await(timeout, timeUnit)可以防止永久等待。 -
初始化计数器的值:应该等于需要等待完成的线程或操作数量。
-
线程安全:CountDownLatch 是线程安全的,可以被多个线程安全地调用。
-
不可重用:CountDownLatch 是一次性工具,计数归零后不能再使用,如果需要循环使用,应该考虑使用 CyclicBarrier。
-
应用场景:
- 主线程等待多个子线程完成
- 模拟高并发场景
- 任务拆分汇总
- 服务启动时等待依赖组件就绪
这就是 CountDownLatch 的常用案例,涵盖了从基础到实战的多个场景,你可以根据自己的需求选择合适的案例进行参考。