本文目录导读:

我将为您提供一个全面的 Exchanger 案例,涵盖基本用法、实际应用场景和注意事项。
Exchanger 基本概念
Exchanger 是 Java 并发包中的同步工具类,用于在两个线程之间交换数据,当两个线程都到达交换点时,它们会交换各自的数据。
案例 1:基本数据交换
import java.util.concurrent.Exchanger;
import java.util.concurrent.TimeUnit;
public class BasicExchangerDemo {
public static void main(String[] args) {
Exchanger<String> exchanger = new Exchanger<>();
// 线程1
new Thread(() -> {
try {
String data1 = "线程1的数据";
System.out.println(Thread.currentThread().getName() +
" 准备交换数据: " + data1);
// 交换数据,等待另一个线程
String received = exchanger.exchange(data1);
System.out.println(Thread.currentThread().getName() +
" 收到数据: " + received);
} catch (InterruptedException e) {
e.printStackTrace();
}
}, "线程1").start();
// 线程2
new Thread(() -> {
try {
// 模拟一些预处理
Thread.sleep(2000);
String data2 = "线程2的数据";
System.out.println(Thread.currentThread().getName() +
" 准备交换数据: " + data2);
// 交换数据
String received = exchanger.exchange(data2);
System.out.println(Thread.currentThread().getName() +
" 收到数据: " + received);
} catch (InterruptedException e) {
e.printStackTrace();
}
}, "线程2").start();
}
}
案例 2:生产消费模式
import java.util.concurrent.Exchanger;
import java.util.ArrayList;
import java.util.List;
import java.util.Random;
public class ProducerConsumerDemo {
// 缓冲区大小
private static final int BUFFER_SIZE = 10;
static class Producer implements Runnable {
private Exchanger<List<Integer>> exchanger;
private List<Integer> buffer;
private int count = 0;
private Random random = new Random();
public Producer(Exchanger<List<Integer>> exchanger) {
this.exchanger = exchanger;
this.buffer = new ArrayList<>(BUFFER_SIZE);
}
@Override
public void run() {
try {
while (count < 20) {
// 填充缓冲区
while (buffer.size() < BUFFER_SIZE && count < 20) {
buffer.add(count++);
System.out.println("生产者生产: " + count);
Thread.sleep(random.nextInt(100));
}
// 缓冲区满了,与消费者交换
System.out.println("生产者缓冲区已满,准备交换: " +
buffer.size() + " 个数据");
buffer = exchanger.exchange(buffer);
System.out.println("生产者收到空缓冲区");
}
// 交换空的缓冲区,通知消费者结束
exchanger.exchange(buffer);
} catch (InterruptedException e) {
e.printStackTrace();
}
}
}
static class Consumer implements Runnable {
private Exchanger<List<Integer>> exchanger;
private List<Integer> buffer;
public Consumer(Exchanger<List<Integer>> exchanger) {
this.exchanger = exchanger;
this.buffer = new ArrayList<>(BUFFER_SIZE);
}
@Override
public void run() {
try {
while (true) {
// 与生产者交换缓冲区
System.out.println("消费者准备交换缓冲区");
buffer = exchanger.exchange(buffer);
System.out.println("消费者收到数据: " + buffer.size() + " 个");
// 如果收到空缓冲区,说明生产者已完成
if (buffer.isEmpty() && !Thread.currentThread().isInterrupted()) {
System.out.println("消费者收到空缓冲区,退出");
break;
}
// 消费数据
for (int data : buffer) {
System.out.println("消费者消费: " + data);
Thread.sleep(50);
}
buffer.clear();
}
} catch (InterruptedException e) {
e.printStackTrace();
}
}
}
public static void main(String[] args) {
Exchanger<List<Integer>> exchanger = new Exchanger<>();
Thread producerThread = new Thread(new Producer(exchanger), "生产者");
Thread consumerThread = new Thread(new Consumer(exchanger), "消费者");
producerThread.start();
consumerThread.start();
// 等待生产者结束
try {
producerThread.join();
Thread.sleep(1000);
consumerThread.interrupt();
} catch (InterruptedException e) {
e.printStackTrace();
}
}
}
案例 3:带超时的交换
import java.util.concurrent.Exchanger;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
public class TimeoutExchangerDemo {
public static void main(String[] args) {
Exchanger<String> exchanger = new Exchanger<>();
// 线程1:正常交换
new Thread(() -> {
try {
String result = exchanger.exchange("线程1的数据",
5, TimeUnit.SECONDS);
System.out.println("线程1收到: " + result);
} catch (InterruptedException e) {
System.err.println("线程1被中断");
} catch (TimeoutException e) {
System.err.println("线程1等待超时");
}
}, "线程1").start();
// 线程2:延迟执行导致线程1超时
new Thread(() -> {
try {
Thread.sleep(3000);
System.out.println("线程2准备交换");
String result = exchanger.exchange("线程2的数据",
5, TimeUnit.SECONDS);
System.out.println("线程2收到: " + result);
} catch (Exception e) {
e.printStackTrace();
}
}, "线程2").start();
}
}
案例 4:综合应用 - 银行转账
import java.util.concurrent.Exchanger;
public class BankTransferDemo {
static class Account {
private String name;
private int balance;
public Account(String name, int balance) {
this.name = name;
this.balance = balance;
}
public int getBalance() {
return balance;
}
public void setBalance(int balance) {
this.balance = balance;
}
@Override
public String toString() {
return name + "账户余额: " + balance;
}
}
public static void main(String[] args) {
// 创建两个账户
Account accountA = new Account("A", 1000);
Account accountB = new Account("B", 500);
System.out.println("初始状态:");
System.out.println(accountA);
System.out.println(accountB);
// 交换余额
Exchanger<Account> exchanger = new Exchanger<>();
// 账户A线程
new Thread(() -> {
try {
System.out.println("账户A准备交换余额");
Account received = exchanger.exchange(accountA);
System.out.println("账户A收到: " + received);
} catch (InterruptedException e) {
e.printStackTrace();
}
}, "账户A线程").start();
// 账户B线程
new Thread(() -> {
try {
System.out.println("账户B准备交换余额");
Account received = exchanger.exchange(accountB);
System.out.println("账户B收到: " + received);
} catch (InterruptedException e) {
e.printStackTrace();
}
}, "账户B线程").start();
}
}
案例 5:任务分配合作
import java.util.concurrent.Exchanger;
import java.util.Arrays;
import java.util.Random;
public class TaskAllocationDemo {
static class Worker {
private String name;
private Exchanger<int[]> exchanger;
private Random random = new Random();
public Worker(String name, Exchanger<int[]> exchanger) {
this.name = name;
this.exchanger = exchanger;
}
public void work() {
try {
// 生成自己要处理的数组
int[] myData = generateData();
System.out.println(name + " 生成数据: " + Arrays.toString(myData));
// 与另一个工人交换数据
int[] receivedData = exchanger.exchange(myData);
System.out.println(name + " 收到数据: " + Arrays.toString(receivedData));
// 处理收到的数据
int sum = Arrays.stream(receivedData).sum();
System.out.println(name + " 处理完成,数据总和: " + sum);
} catch (InterruptedException e) {
e.printStackTrace();
}
}
private int[] generateData() {
int[] data = new int[5];
for (int i = 0; i < data.length; i++) {
data[i] = random.nextInt(100);
}
return data;
}
}
public static void main(String[] args) {
Exchanger<int[]> exchanger = new Exchanger<>();
Worker worker1 = new Worker("工人1", exchanger);
Worker worker2 = new Worker("工人2", exchanger);
new Thread(worker1::work).start();
new Thread(worker2::work).start();
}
}
使用注意事项
交换匹配规则
- 必须是成对的线程进行交换
- 如果第三者插入交换,会导致行为不可预测
超时处理
// 设置超时时间,避免无限等待
try {
String result = exchanger.exchange(data, 10, TimeUnit.SECONDS);
} catch (TimeoutException e) {
System.err.println("等待超时");
}
中断处理
try {
String result = exchanger.exchange(data);
} catch (InterruptedException e) {
Thread.currentThread().interrupt(); // 保持中断状态
System.err.println("交换被中断");
}
空值判定
// 可以交换null值,但要注意后续处理
Object received = exchanger.exchange(null);
if (received != null) {
// 处理数据
}
| 场景 | 推荐用法 | 注意事项 |
|---|---|---|
| 简单数据交换 | 直接使用 exchange() |
确保两个线程同时到达 |
| 生产消费模式 | 使用 Exchange 交换缓冲区 | 注意缓冲区大小管理 |
| 需要超时控制 | 使用 exchange(data, timeout, unit) |
处理 TimeoutException |
| 需要中断控制 | 捕获 InterruptedException |
正确处理线程中断 |
| 防止死锁 | 设置合理的超时时间 | 避免无限等待 |
性能考虑
// 高并发场景下的性能优化
public class PerformanceExchangerDemo {
// 1. 预估并发量,预创建Exchanger
private static final int CONCURRENCY_LEVEL = 8;
private static final Exchanger<byte[]>[] exchangers = new Exchanger[CONCURRENCY_LEVEL];
static {
for (int i = 0; i < CONCURRENCY_LEVEL; i++) {
exchangers[i] = new Exchanger<>();
}
}
// 2. 使用ThreadLocal减少竞争
private static final ThreadLocal<Exchanger<byte[]>> localExchanger =
ThreadLocal.withInitial(() -> new Exchanger<>());
public static void main(String[] args) {
// 示例代码
byte[] data = new byte[1024];
try {
byte[] received = localExchanger.get().exchange(data);
// 处理数据
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
}
Exchanger 特别适合解决两个线程需要交换数据的场景,如生产者-消费者模型、遗传算法、流水线作业等,使用时要特别注意线程配对和超时处理,以防出现死锁或无限等待的情况。