Exchanger案例

wen java案例 2

本文目录导读:

Exchanger案例

  1. Exchanger 基本概念
  2. 案例 1:基本数据交换
  3. 案例 2:生产消费模式
  4. 案例 3:带超时的交换
  5. 案例 4:综合应用 - 银行转账
  6. 案例 5:任务分配合作
  7. 使用注意事项
  8. 性能考虑

我将为您提供一个全面的 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 特别适合解决两个线程需要交换数据的场景,如生产者-消费者模型、遗传算法、流水线作业等,使用时要特别注意线程配对和超时处理,以防出现死锁或无限等待的情况。

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