并发工具类详解
飞翔科技项目作战室,周五下午。
朱璐(测试组长,满头大汗):"完了完了,集成测试又挂了!10个微服务模块,前置检查模块还没跑完,业务模块就开始请求数据库了,直接雪崩。"
小崔(初级开发):"那让它们按顺序执行不就行了?"
白歌(首席架构师):"按顺序执行效率太低。应该用 CountDownLatch——让业务模块等前置检查全部就绪后再统一启动。就像起跑线上的运动员,等发令枪一响同时出发。"
大翔(CTO):"对,还有压测时的并发控制——我们上次没用 Semaphore,500 个并发直接把下游服务打挂了。"
Frank(后端专家):"Exactly. JDK 8 provides a whole toolkit: CountDownLatch for one-time coordination, CyclicBarrier for reusable barriers, Semaphore for resource limiting, and even Exchanger for two-party data swap. Each solves a different concurrency coordination problem."
小崔:"这么多工具类,到底什么时候用哪个啊?"
白歌:"问得好。今晚我们就来逐个击破。"
一、定义表
| 工具类 | 定义 | JDK 8 所在包 | 核心比喻 |
|---|---|---|---|
| CountDownLatch | 一次性门闩:主线程等待一组工作线程全部就绪后继续执行 | java.util.concurrent | 起跑发令枪——等所有人到齐才能开枪 |
| CyclicBarrier | 可重用屏障:一组线程互相等待,全部到达屏障后一起继续 | java.util.concurrent | 团队打卡——人到齐了一起出发,可循环多轮 |
| Semaphore | 信号量:控制同时访问资源的线程数 | java.util.concurrent | 停车场——N个车位,满了就得等 |
| Exchanger | 交换器:两个线程在汇合点交换数据 | java.util.concurrent | 双人交接——你把你那份给我,我把我这份给你 |
| Phaser | 分层阶段器:比 CyclicBarrier 更灵活的多阶段屏障(JDK 7引入,JDK 8增强) | java.util.concurrent | 接力赛——每一棒完成后大家同步,再进入下一棒 |
二、CountDownLatch —— 一次性门闩
2.1 工作原理
主线程: await() ────────────阻塞────────────→ 继续执行
↑ 计数归零 ↑
工作线程1: countDown() ──┐
工作线程2: countDown() ──┤
工作线程3: countDown() ──┘
- 计数器的值在构造时设定,只能减少不能增加。
countDown()将计数器减 1。await()阻塞直到计数器归零。- 一次性:计数器归零后无法重置,需要新实例。
2.2 Mermaid 流程图
三、CyclicBarrier —— 可重用屏障
3.1 CountDownLatch vs CyclicBarrier 对比
| 对比维度 | CountDownLatch | CyclicBarrier |
|---|---|---|
| 计数方向 | 递减(只能减不能增) | 递增(到达线程数累加) |
| 可重用性 | 不可重用,归零后失效 | 可重用,屏障打开后自动重置 |
| 等待方 | 通常主线程等待工作线程 | 工作线程互相等待 |
| 回调 | 无 | 支持 Runnable barrierAction |
| 破坏机制 | 无 | reset() 可强制破坏屏障 |
四、Semaphore —— 信号量(令牌桶)
五、代码示例
示例一:飞翔科技多模块启动同步——CountDownLatch
场景:飞翔科技电商平台启动时需要先完成 5 个前置模块的初始化(数据库连接池、Redis 缓存、消息队列、配置中心、日志系统),全部就绪后主服务才对外暴露 HTTP 端口。
import java.util.concurrent.*;
/**
* 飞翔科技 —— 服务启动多模块同步(CountDownLatch)
*/
public class ServiceStartupLatch {
// 模拟模块初始化
static class ModuleInitializer implements Runnable {
private final String moduleName;
private final int initTimeMs;
private final CountDownLatch latch;
public ModuleInitializer(String moduleName, int initTimeMs, CountDownLatch latch) {
this.moduleName = moduleName;
this.initTimeMs = initTimeMs;
this.latch = latch;
}
@Override
public void run() {
try {
System.out.printf("[%s] %s 开始初始化...%n",
Thread.currentThread().getName(), moduleName);
Thread.sleep(initTimeMs); // 模拟初始化耗时
System.out.printf("[%s] %s 初始化完成 ✓%n",
Thread.currentThread().getName(), moduleName);
} catch (InterruptedException e) {
System.err.println(moduleName + " 初始化被中断!");
} finally {
latch.countDown(); // 关键:无论成功/失败都要倒计数,防止死锁
}
}
}
public static void main(String[] args) throws InterruptedException {
int moduleCount = 5;
CountDownLatch latch = new CountDownLatch(moduleCount);
ExecutorService executor = Executors.newFixedThreadPool(moduleCount);
System.out.println("===== 飞翔科技电商平台启动中... =====\n");
long startTime = System.currentTimeMillis();
// 提交 5 个模块初始化任务
executor.execute(new ModuleInitializer("数据库连接池", 2000, latch));
executor.execute(new ModuleInitializer("Redis 缓存集群", 1500, latch));
executor.execute(new ModuleInitializer("RocketMQ 消息队列", 1800, latch));
executor.execute(new ModuleInitializer("Nacos 配置中心", 1000, latch));
executor.execute(new ModuleInitializer("ELK 日志系统", 1200, latch));
System.out.println("主线程: 等待所有模块初始化完成...\n");
// 主线程阻塞,直到所有模块初始化完成
latch.await();
long elapsed = System.currentTimeMillis() - startTime;
System.out.println("\n==========================================");
System.out.printf("所有 %d 个模块初始化完成!总耗时: %d ms%n", moduleCount, elapsed);
System.out.println("主服务开始监听端口 8080...");
System.out.println("飞翔科技电商平台启动成功!");
System.out.println("==========================================");
executor.shutdown();
}
}
控制台输出:
===== 飞翔科技电商平台启动中... =====
[pool-1-thread-1] 数据库连接池 开始初始化...
[pool-1-thread-2] Redis 缓存集群 开始初始化...
[pool-1-thread-3] RocketMQ 消息队列 开始初始化...
[pool-1-thread-4] Nacos 配置中心 开始初始化...
[pool-1-thread-5] ELK 日志系统 开始初始化...
主线程: 等待所有模块初始化完成...
[pool-1-thread-4] Nacos 配置中心 初始化完成 ✓
[pool-1-thread-5] ELK 日志系统 初始化完成 ✓
[pool-1-thread-2] Redis 缓存集群 初始化完成 ✓
[pool-1-thread-3] RocketMQ 消息队列 初始化完成 ✓
[pool-1-thread-1] 数据库连接池 初始化完成 ✓
==========================================
所有 5 个模块初始化完成!总耗时: 2018 ms
主服务开始监听端口 8080...
飞翔科技电商平台启动成功!
==========================================
示例二:压力测试并发限制——Semaphore + CyclicBarrier 联合
场景:朱璐的测试团队需要对订单接口进行压测。Semaphore 限制同时打向接口的并发数为 10,CyclicBarrier 让所有虚拟用户准备好后同时发起请求。
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;
/**
* 飞翔科技 —— 订单接口压力测试(Semaphore + CyclicBarrier)
*/
public class OrderApiLoadTest {
// 模拟订单接口
static class OrderService {
private final AtomicInteger successCount = new AtomicInteger(0);
private final AtomicInteger failCount = new AtomicInteger(0);
public String createOrder(int userId, String product) {
// 模拟接口处理耗时 50~150ms
long processTime = 50 + (long)(Math.random() * 100);
try {
Thread.sleep(processTime);
successCount.incrementAndGet();
return String.format("用户%d购买%s成功(耗时%dms)", userId, product, processTime);
} catch (InterruptedException e) {
failCount.incrementAndGet();
return "订单创建失败";
}
}
public void printStats() {
System.out.printf("\n订单统计: 成功=%d, 失败=%d%n",
successCount.get(), failCount.get());
}
}
public static void main(String[] args) throws Exception {
final int TOTAL_USERS = 50; // 总虚拟用户数
final int MAX_CONCURRENT = 10; // 最大并发数(Semaphore 许可数)
Semaphore semaphore = new Semaphore(MAX_CONCURRENT);
// CyclicBarrier: 等 50 个线程都就绪后一起开始
CyclicBarrier barrier = new CyclicBarrier(TOTAL_USERS,
() -> System.out.println("\n🚀 所有虚拟用户已就绪,开始压测!\n"));
OrderService orderService = new OrderService();
ExecutorService executor = Executors.newFixedThreadPool(TOTAL_USERS);
CountDownLatch doneLatch = new CountDownLatch(TOTAL_USERS);
String[] products = {"iPhone 15", "MacBook Pro", "AirPods", "iPad",
"Apple Watch", "HomePod"};
System.out.println("===== 飞翔科技订单接口压力测试 =====");
System.out.printf("虚拟用户数: %d | 最大并发: %d%n", TOTAL_USERS, MAX_CONCURRENT);
System.out.println("等待所有用户就绪...");
long startTime = System.currentTimeMillis();
for (int i = 0; i < TOTAL_USERS; i++) {
final int userId = 1000 + i;
final String product = products[i % products.length];
executor.execute(() -> {
try {
// 第1步:等待所有用户就绪(CyclicBarrier)
barrier.await();
// 第2步:获取信号量许可(控制并发)
semaphore.acquire();
try {
String result = orderService.createOrder(userId, product);
System.out.printf("[%s] %s%n",
Thread.currentThread().getName(), result);
} finally {
semaphore.release(); // 关键:finally 中释放
}
} catch (InterruptedException | BrokenBarrierException e) {
Thread.currentThread().interrupt();
} finally {
doneLatch.countDown();
}
});
}
// 等待所有用户完成
doneLatch.await();
long elapsed = System.currentTimeMillis() - startTime;
orderService.printStats();
System.out.printf("压测完成!总耗时: %d ms | 平均TPS: %.0f%n",
elapsed, TOTAL_USERS * 1000.0 / elapsed);
executor.shutdown();
}
}
控制台输出(部分):
===== 飞翔科技订单接口压力测试 =====
虚拟用户数: 50 | 最大并发: 10
等待所有用户就绪...
🚀 所有虚拟用户已就绪,开始压测!
[pool-1-thread-1] 用户1000购买iPhone 15成功(耗时78ms)
[pool-1-thread-3] 用户1002购买AirPods成功(耗时92ms)
[pool-1-thread-5] 用户1004购买Apple Watch成功(耗时65ms)
...
[pool-1-thread-48] 用户1047购买HomePod成功(耗时112ms)
[pool-1-thread-50] 用户1049购买iPhone 15成功(耗时133ms)
订单统计: 成功=50, 失败=0
压测完成!总耗时: 634 ms | 平均TPS: 79
示例三:Exchanger —— 双线程数据交换
场景:飞翔科技的数据清洗流水线中,两个 worker 线程分别处理原始日志和格式化日志,处理完成后在汇合点交换数据,互相校验。
import java.util.concurrent.Exchanger;
/**
* 飞翔科技 —— 双工数据校验交换(Exchanger)
*/
public class LogDataExchanger {
public static void main(String[] args) {
Exchanger<String> exchanger = new Exchanger<>();
// 线程A:处理原始访问日志
new Thread(() -> {
try {
String rawLog = "192.168.1.100|/api/order|2025-01-15 10:23:45|200|32ms";
System.out.println("[原始日志线程] 处理完成: " + rawLog);
// 模拟处理耗时
Thread.sleep(1000);
// 交换数据
System.out.println("[原始日志线程] 等待与格式化线程交换数据...");
String formattedLog = exchanger.exchange(rawLog);
System.out.println("[原始日志线程] 收到格式化日志: " + formattedLog);
// 校验数据一致性
if (formattedLog.contains("200") && formattedLog.contains("/api/order")) {
System.out.println("[原始日志线程] 数据校验通过 ✓");
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}, "raw-log-worker").start();
// 线程B:处理格式化日志
new Thread(() -> {
try {
String formattedLog = "[INFO] 192.168.1.100 - \"GET /api/order HTTP/1.1\" 200 32ms";
System.out.println("[格式化线程] 处理完成: " + formattedLog);
// 模拟处理耗时
Thread.sleep(2000);
// 交换数据
System.out.println("[格式化线程] 等待与原始日志线程交换数据...");
String rawLog = exchanger.exchange(formattedLog);
System.out.println("[格式化线程] 收到原始日志: " + rawLog);
// 反向校验
if (rawLog.contains("192.168.1.100") && rawLog.contains("32ms")) {
System.out.println("[格式化线程] 数据校验通过 ✓");
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}, "formatted-log-worker").start();
}
}
控制台输出:
[原始日志线程] 处理完成: 192.168.1.100|/api/order|2025-01-15 10:23:45|200|32ms
[格式化线程] 处理完成: [INFO] 192.168.1.100 - "GET /api/order HTTP/1.1" 200 32ms
[原始日志线程] 等待与格式化线程交换数据...
[格式化线程] 等待与原始日志线程交换数据...
[原始日志线程] 收到格式化日志: [INFO] 192.168.1.100 - "GET /api/order HTTP/1.1" 200 32ms
[原始日志线程] 数据校验通过 ✓
[格式化线程] 收到原始日志: 192.168.1.100|/api/order|2025-01-15 10:23:45|200|32ms
[格式化线程] 数据校验通过 ✓
六、Phaser —— 分层阶段同步(JDK 7+)
Phaser 是 CyclicBarrier 和 CountDownLatch 的超级合体,支持动态注册/注销参与者、多阶段同步。
import java.util.concurrent.Phaser;
/**
* 飞翔科技 —— 多阶段发布流程(Phaser)
*/
public class ReleasePipelinePhaser {
static class ReleasePhase implements Runnable {
private final String phaseName;
private final int durationMs;
private final Phaser phaser;
public ReleasePhase(String phaseName, int durationMs, Phaser phaser) {
this.phaseName = phaseName;
this.durationMs = durationMs;
this.phaser = phaser;
}
@Override
public void run() {
// 阶段1:代码检查
doPhase("Phase 1 - 代码静态检查");
// 阶段2:单元测试
doPhase("Phase 2 - 单元测试");
// 阶段3:打包部署
doPhase("Phase 3 - 打包部署");
}
private void doPhase(String phaseDesc) {
System.out.printf("[%s] %s - 阶段%d 开始%n",
Thread.currentThread().getName(),
phaseName, phaser.getPhase());
try {
Thread.sleep(durationMs);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
System.out.printf("[%s] %s - 阶段%d 完成%n",
Thread.currentThread().getName(),
phaseName, phaser.getPhase());
phaser.arriveAndAwaitAdvance(); // 到达并等待其他参与者
}
}
public static void main(String[] args) {
Phaser phaser = new Phaser() {
@Override
protected boolean onAdvance(int phase, int registeredParties) {
System.out.println("===== 阶段 " + phase + " 全部完成,进入下一阶段 =====\n");
// 返回 true 则终止 phaser
return phase >= 2; // 3个阶段(0,1,2)后终止
}
};
// 注册3个参与者
phaser.register();
phaser.register();
phaser.register();
System.out.println("===== 飞翔科技发布流水线启动 (Phaser) =====\n");
new Thread(new ReleasePhase("订单服务", 500, phaser), "order-svc").start();
new Thread(new ReleasePhase("用户服务", 600, phaser), "user-svc").start();
new Thread(new ReleasePhase("支付服务", 400, phaser), "pay-svc").start();
// 主线程也参与同步
for (int i = 0; i < 3; i++) {
phaser.arriveAndAwaitAdvance();
}
System.out.println("所有阶段完成,Phaser 终止。");
}
}
七、易错场景
7.1 CountDownLatch 忘记在 finally 中 countDown 导致死锁
// ❌ 错误:异常时未 countDown,主线程永久阻塞
CountDownLatch latch = new CountDownLatch(3);
for (int i = 0; i < 3; i++) {
new Thread(() -> {
int result = 1 / 0; // 抛出异常
latch.countDown(); // 这行永远不会执行!
}).start();
}
latch.await(); // 死锁!永远等不到计数归零
System.out.println("这行永远不会输出");
正确做法:始终在 finally 块中调用 countDown()。
7.2 CyclicBarrier 的 barrierAction 中抛出异常
// ❌ 危险:barrierAction 抛异常会破坏屏障,其他线程收到 BrokenBarrierException
CyclicBarrier barrier = new CyclicBarrier(3, () -> {
throw new RuntimeException("barrierAction 异常"); // 导致屏障破坏
});
// 所有 await() 的线程都会收到 BrokenBarrierException
for (int i = 0; i < 3; i++) {
new Thread(() -> {
try {
barrier.await();
} catch (BrokenBarrierException e) {
System.out.println("屏障已被破坏: " + e.getMessage());
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}).start();
}
正确做法:barrierAction 必须做好异常保护,屏障破坏后调用 reset() 重置。
7.3 Semaphore 获取许可后未在 finally 中 release
// ❌ 危险:异常导致许可永不释放 → 其他线程饿死
Semaphore sem = new Semaphore(2);
sem.acquire();
try {
int result = 1 / 0; // 异常!
} // 缺少 finally { sem.release(); }
// 许可泄漏!永久少了一个许可
正确做法:acquire() 和 release() 必须成对出现在 try-finally 中。
7.4 Exchanger 只有一个线程到达会导致永久阻塞
// ❌ 错误:如果只有一个线程调用 exchange(),该线程永久阻塞
Exchanger<String> ex = new Exchanger<>();
new Thread(() -> {
try {
ex.exchange("data"); // 永久阻塞,因为没有第二个线程来交换
System.out.println("永远不会输出");
} catch (InterruptedException e) {}
}).start();
正确做法:使用 exchange(V x, long timeout, TimeUnit unit) 设置超时。
八、面试考点
Q1:CountDownLatch 和 CyclicBarrier 的区别?(高频)
答:
| CountDownLatch | CyclicBarrier |
|---|---|
| 计数器只能减,不能重置 | 计数器自动重置,可循环使用 |
| 通常是一个线程等待多个线程完成 | 多个线程互相等待,全部到达后同时继续 |
countDown() 后线程继续执行不阻塞 | await() 后线程阻塞直到所有线程都到达 |
| 无回调 | 支持 barrierAction,屏障打开前执行 |
Q2:Semaphore 的公平模式和非公平模式有什么区别?
答:Semaphore(int permits, boolean fair) 的 fair 参数控制是否启用公平模式。
- 非公平模式(默认):新来的线程可以"插队"直接抢许可,吞吐量更高,但可能导致某些线程饥饿。
- 公平模式:按 FIFO 顺序分配许可,严格排队,公平性更好但吞吐量较低(涉及挂起/唤醒开销)。
- JDK 8 中 Semaphore 基于 AQS(AbstractQueuedSynchronizer)实现。
Q3:CyclicBarrier 的 broken 状态是怎么回事?
答:以下情况会导致屏障进入 broken 状态:
barrierAction抛出未捕获的异常。- 某个等待线程被中断。
- 某个等待线程超时(使用
await(timeout))。 - 主动调用
reset()方法。
broken 状态下,所有等待中的线程会收到 BrokenBarrierException。必须调用 reset() 才能使屏障恢复正常。
Q4:Exchanger 的底层实现原理?(追问)
答:Exchanger 内部使用 slot exchange 机制。两个线程在交换点时,第一个到达的线程将数据写入 Node 的 item 字段并自旋等待;第二个到达的线程读取该数据并写入自己的数据,唤醒第一个线程。JDK 8 中使用 sun.misc.Unsafe 的 CAS 操作实现无锁交换,并加入了自旋优化和 park/unpark 机制。
九、选型速查表
| 需求 | 推荐工具 |
|---|---|
| 等 N 个子任务全部完成再继续 | CountDownLatch |
| N 个线程互相等待,全部就绪后同时出发 | CyclicBarrier |
| 限制同时访问某资源的线程数 | Semaphore |
| 两个线程交换数据 | Exchanger |
| 多阶段、动态参与者同步 | Phaser |
| 需要可重复使用的屏障 + 回调 | CyclicBarrier |
| 一次性等待 + 不需要回调 | CountDownLatch |