1. CountDownLatch:按信号启动
在多线程世界中,经常需要让一组线程协同工作——让它们一起开始、一起结束,或一起进入下一阶段。比如:
想象一场比赛。赛车都在起跑线上——有的已经热好发动机,有的还在检查轮胎。但在裁判挥旗之前,谁也不会起步。这就是协调的任务。
或者另一个例子:你和朋友一起做晚饭——有人切菜,有人烧水,有人在找盐去哪儿了。重点是,在开始正式烹饪之前,所有人都要把准备工作完成。
针对这类场景,Java 给出了现成的同步工具——安全、易懂,而且不用承受 wait() 和 notify() 的痛苦。其中一个最实用的就是 CountDownLatch。它像一个计数锁:在计数归零之前,“门”是关着的,谁也不能继续往下执行;当所有人都报到后——latch 打开,线程同步冲向下一步。
CountDownLatch
CountDownLatch 是一个“一次性阀门”,它允许一个或多个线程等待,直到其他线程完成了指定数量的操作。
这就像马拉松的起跑:所有跑者站在起跑线,等待发令枪。一旦裁判开枪(计数减到 0)——大家一起开跑。
到底怎么工作的
CountDownLatch 就像线程的起跑哨。创建时你设定一个数字——比如 3。这就像是比赛开始前需要收到的三个信号。
那些需要等待起跑的线程调用 await()。它们已在起跑线整装待发,但暂时还踩着刹车。其他线程在完成准备的过程中,会按时调用 countDown()——相当于发出信号:“我准备好了!”。
当计数器归零——砰!——所有等待的线程同时起跑。
但请记住:CountDownLatch 是一次性的。计数器归零后就无法再复原。它不是左轮手枪,而是鞭炮:响完就没了。
示例:等待 N 个任务完成
import java.util.concurrent.CountDownLatch;
public class LatchDemo {
public static void main(String[] args) throws InterruptedException {
int workers = 3;
CountDownLatch latch = new CountDownLatch(workers);
for (int i = 1; i <= workers; i++) {
int id = i;
new Thread(() -> {
System.out.println("工作者 " + id + " 开始工作");
try { Thread.sleep(500 + id * 200); } catch (InterruptedException ignored) {}
System.out.println("工作者 " + id + " 完成工作");
latch.countDown(); // 递减计数器
}).start();
}
System.out.println("主线程正在等待所有工作者完成...");
latch.await(); // 等待所有工作者结束
System.out.println("所有工作者都结束了!继续主流程。");
}
}
输出:
主线程正在等待所有工作者完成...
工作者 1 开始工作
工作者 2 开始工作
工作者 3 开始工作
工作者 1 完成工作
工作者 2 完成工作
工作者 3 完成工作
所有工作者都结束了!继续主流程。
示例:按“信号”同时启动
CountDownLatch startSignal = new CountDownLatch(1);
for (int i = 0; i < 5; i++) {
new Thread(() -> {
try {
System.out.println(Thread.currentThread().getName() + " 在等待开始");
startSignal.await(); // 等待信号
System.out.println(Thread.currentThread().getName() + " 开始!");
} catch (InterruptedException ignored) {}
}).start();
}
Thread.sleep(1000);
System.out.println("开始信号!");
startSignal.countDown(); // 所有线程同时启动
2. CyclicBarrier:多次阶段、屏障动作
CyclicBarrier:在“篝火处”会合
CyclicBarrier 是线程的会合点。每个线程沿着自己的路线各自做事,然后大家在“屏障”处集合——好比在山中的篝火旁。等人齐了,屏障打开,团队一起继续前进。
和 CountDownLatch 的主要区别在于——这个屏障可以反复使用。每次集体停靠之后它都会“重置”,小队可以继续奔向下一个阶段。
想象一下:一支登山队走很长的路线。每个人的节奏不同:有人拍蝴蝶,有人找 Wi‑Fi。但在每个山口他们都会在篝火旁会合,互相等待,然后决定接下来走向哪里。这就是 CyclicBarrier 的运作方式。
如何工作
你创建屏障并指定需要集合的参与者数量,比如 4。每个线程到达检查点后调用 await()——然后等待其他人。当四个人都到齐,屏障“咔哒”一声放行,所有人一起继续。
你甚至可以设置“屏障动作”——当小队集合时,只执行一次的一段代码。比如点燃那堆篝火,或记录日志:“阶段完成,继续前进”。为此可在构造函数中传入一个 Runnable。
重要:不同于一次性的 CountDownLatch,CyclicBarrier 是可重复使用的。每次集合后它都会再次待命——就像一堆随时可以再次点燃的“旅行篝火”。
示例:阶段同步
import java.util.concurrent.CyclicBarrier;
public class BarrierDemo {
public static void main(String[] args) {
int parties = 3;
CyclicBarrier barrier = new CyclicBarrier(parties, () -> {
System.out.println("所有人都到达屏障!开始新阶段。");
});
for (int i = 1; i <= parties; i++) {
int id = i;
new Thread(() -> {
try {
System.out.println("线程 " + id + " 在阶段 1 中工作");
Thread.sleep(300 + id * 200);
System.out.println("线程 " + id + " 等待屏障");
barrier.await(); // 等待其他人
System.out.println("线程 " + id + " 在阶段 2 中工作");
Thread.sleep(200 + id * 100);
System.out.println("线程 " + id + " 等待屏障(2)");
barrier.await(); // 再次等待
System.out.println("线程 " + id + " 已完成工作");
} catch (Exception e) {
System.out.println("错误:" + e);
}
}).start();
}
}
}
输出:
线程 1 在阶段 1 中工作
线程 2 在阶段 1 中工作
线程 3 在阶段 1 中工作
线程 1 等待屏障
线程 2 等待屏障
线程 3 等待屏障
所有人都到达屏障!开始新阶段。
线程 1 在阶段 2 中工作
...
屏障动作
可以在 CyclicBarrier 的构造函数中传入一个动作(Runnable),当所有线程到达屏障时,该动作将执行一次(例如更新状态、输出日志)。
陷阱:如果某个线程挂了怎么办?
如果某个线程抛出异常或没到达屏障,其他线程会一直等待——或者抛出 BrokenBarrierException。此时屏障“损坏”,需要重新创建。
下面这个部分,我们把它写得更生动一些,让它自然承接“乐团”的比喻:
3. Phaser:大型音乐会的能干指挥
Phaser 有点像“超级屏障”。它结合了 CountDownLatch 和 CyclicBarrier 的优点,但更灵活。就像一支乐团,乐手可以在不同乐章之间进进出出,而指挥依然能确保每个乐章都在大家准备好时再开始。
与普通屏障不同,Phaser 能按阶段工作——阶段一个接一个地推进。有人只在第一部分演奏,有人后来加入,也有人提前离场——这些对 Phaser 来说都不成问题。
如何工作
先创建 Phaser,通常会指定参与者数量——parties。每个线程先注册(register()),执行自己的“乐段”,在阶段末尾调用 arriveAndAwaitAdvance()——表示已完成并等待其他人。当所有人都到达该点,Phaser 切换到下一阶段,流程重复。
如果不再需要某个参与者——它可以优雅地“鞠躬”退场,通过 arriveAndDeregister()。新的参与者也可以在演出过程中加入——通过 register()。
什么时候 Phaser 比 Barrier 更好
Phaser 更适用于你的程序不只一个节奏,而是多个节奏的情况:
- 线程数量会在运行时变化,
- 有多个阶段,且不是所有参与者都必须参与所有阶段,
- 或者你想要最大灵活性而不想手写复杂的同步逻辑。
本质上,Phaser 是“指挥”。它不仅挥动指挥棒,还能适配乐团的编制、乐章的数量,甚至应对有人迟到或提前离场。
示例:动态线程数的分阶段处理
import java.util.concurrent.Phaser;
public class PhaserDemo {
public static void main(String[] args) {
Phaser phaser = new Phaser(1); // 主线程
for (int i = 1; i <= 3; i++) {
phaser.register(); // 注册参与者
int id = i;
new Thread(() -> {
for (int phase = 1; phase <= 2; phase++) {
System.out.println("线程 " + id + " 在阶段 " + phase + " 中工作");
try { Thread.sleep(200 + id * 100); } catch (InterruptedException ignored) {}
phaser.arriveAndAwaitAdvance(); // 等待其他人
}
System.out.println("线程 " + id + " 已完成工作");
phaser.arriveAndDeregister(); // 从 phaser 注销
}).start();
}
// 主线程也参与各阶段
for (int phase = 1; phase <= 2; phase++) {
phaser.arriveAndAwaitAdvance();
System.out.println("主线程:阶段 " + phase + " 已完成");
}
phaser.arriveAndDeregister();
System.out.println("所有阶段都已完成!");
}
}
要点:
- 可以在运行时添加/移除参与者。
- 可以获取当前阶段编号:phaser.getPhase()。
- 可以结束 phaser:phaser.forceTermination()。
4. Exchanger:线程间的数据成对交换
Exchanger<T> 是用于两个线程之间交换数据的同步器。每个线程调用 exchange(data),当两者会合时,它们交换各自的数据。
类比:两名快递员在十字路口会面并交换包裹。
如何工作?
- 一个线程调用 exchange(data1)——等待第二个线程。
- 第二个线程调用 exchange(data2)——双方各自拿到对方的数据。
- 如果第二个线程没到——第一个线程会等待(可以设置超时)。
示例:producer 与 consumer 的缓冲区交换
import java.util.concurrent.Exchanger;
public class ExchangerDemo {
public static void main(String[] args) {
Exchanger<String> exchanger = new Exchanger<>();
// Producer
new Thread(() -> {
String data = "来自 producer 的数据";
try {
System.out.println("Producer:正在发送数据");
String response = exchanger.exchange(data);
System.out.println("Producer:收到响应:" + response);
} catch (InterruptedException e) {
e.printStackTrace();
}
}).start();
// Consumer
new Thread(() -> {
try {
String received = exchanger.exchange("来自 consumer 的回复");
System.out.println("Consumer:收到数据:" + received);
} catch (InterruptedException e) {
e.printStackTrace();
}
}).start();
}
}
输出:
Producer:正在发送数据
Consumer:收到数据:来自 producer 的数据
Producer:收到响应:来自 consumer 的回复
应用场景:
- 在线程之间交换缓冲区(例如,一个从文件读取,另一个写入网络)。
- 两个线程之间的阶段同步。
5. 实战:并行流水线处理
任务:“游戏 tick”(阶段)
假设我们有多个线程,每个线程负责游戏世界的一部分(例如物理、AI、渲染)。所有线程需要在每个“tick”(阶段)上进行同步,以避免不同步。
解决方案:使用 CyclicBarrier 或 Phaser。
import java.util.concurrent.CyclicBarrier;
public class GameTickDemo {
public static void main(String[] args) {
int subsystems = 3;
CyclicBarrier barrier = new CylicBarrier(subsystems, () -> {
System.out.println("所有子系统完成了本次 tick。开始下一次。");
});
for (int i = 1; i <= subsystems; i++) {
int id = i;
new Thread(() -> {
for (int tick = 1; tick <= 5; tick++) {
System.out.println("子系统 " + id + " 在第 " + tick + " 次 tick 中工作");
try { Thread.sleep(100 + id * 50); } catch (InterruptedException ignored) {}
try {
barrier.await();
} catch (Exception e) {
e.printStackTrace();
}
}
}).start();
}
}
}
任务:大量 worker 的“阀门”式统一起跑
假设我们有 100 个 worker 线程,它们在完成准备后需要同时启动(例如做压力测试)。
解决方案:使用 CountDownLatch。
import java.util.concurrent.CountDownLatch;
public class MassStartDemo {
public static void main(String[] args) throws InterruptedException {
int workers = 100;
CountDownLatch ready = new CountDownLatch(workers);
CountDownLatch start = new CountDownLatch(1);
for (int i = 0; i < workers; i++) {
new Thread(() -> {
System.out.println("线程已准备好");
ready.countDown(); // 发出就绪信号
try {
start.await(); // 等待统一信号
System.out.println("线程开始!");
} catch (InterruptedException ignored) {}
}).start();
}
ready.await(); // 等待所有线程准备就绪
System.out.println("所有人都准备就绪!开始!");
start.countDown(); // 发出启动信号
}
}
6. 使用同步器时的常见错误
错误 1:把 CountDownLatch 当成可复用的屏障。
CountDownLatch 是一次性的!计数归零后不能“重新装填”。可复用的阶段请用 CyclicBarrier 或 Phaser。
错误 2:没有处理异常(InterruptedException、BrokenBarrierException)。
await() 可能会抛出异常——务必处理,否则线程可能“挂起”或以错误结束。特别留意 InterruptedException 和 BrokenBarrierException。
错误 3:某个线程没到达屏障。
如果有线程“挂了”或没有调用 await(),其他线程会一直等待(或抛出 BrokenBarrierException)。务必确保所有参与者都能到达屏障。
错误 4:在 Phaser 中忘记调用 deregister()。
如果线程已经结束却没有调用 arriveAndDeregister(),Phaser 会一直等待这个“僵尸”参与者。务必正确地将线程从 Phaser 中移除。
错误 5:把 Exchanger 用在两个以上的线程之间。
Exchanger 只能用于两个线程之间的数据交换。线程数更多会导致死锁。
错误 6:在不了解工作方式的情况下混用不同同步器。
不要为同一组线程同时使用多个不同的屏障/闩锁——这容易造成混乱和卡死。
GO TO FULL VERSION