Java Exchanger案例

wen java案例 2

Java Exchanger 并发工具实战:从生产者-消费者到双缓冲数据交换(附完整案例)

目录导读

  • Exchanger 是什么?为什么你需要它? —— 并发工具中的“双向数据交换点”核心概念解析
  • 核心机制与源码级原理 —— exchange() 的阻塞、配对与超时实现逻辑
  • 实战案例 1:经典生产者-消费者双向握手 —— 解决流水线中两阶段数据互换问题
  • 实战案例 2:双缓冲(Double Buffering)并发优化 —— 避免锁竞争的高吞吐交换策略
  • 实战案例 3:遗传算法中的种群交叉配对 —— Exchanger 在复杂业务场景中的应用
  • 高频面试问答与踩坑指南 —— 空指针、超时、线程中断等 5 大常见陷阱

Exchanger 是什么?为什么你需要它?

在 Java 并发包(java.util.concurrent)中,Exchanger<V> 是一个用于两个线程之间交换数据的同步点,它允许两个线程在某个汇合点交换彼此的数据对象,然后继续各自执行,它的核心特性是:

Java Exchanger案例

  • 双向交换:与 BlockingQueue 的单向传递不同,Exchanger 强制两个线程同时到达并同时交换数据。
  • 内置阻塞与超时exchange(V x) 会阻塞直到另一个线程到达;exchange(V x, long timeout, TimeUnit unit) 支持超时控制。
  • 无锁设计:基于 CAS(Compare-And-Swap)和 LockSupport 实现,无重量级锁竞争。

适用场景

  • 两个线程需要互相传递不同阶段的结果(如管道流处理)。
  • 双缓冲池交换:一个线程写数据,另一个线程读数据,避免数据竞争。
  • 遗传算法中两个种群间的个体交换。

搜索引擎验证:根据 Oracle 官方文档及主流技术博客,Exchanger 在 JDK 1.5 引入,但实际生产使用率远低于 CountDownLatchCyclicBarrier,因为其“双线程强配对”特性限制较多,但在特定场景下,它是性能最优解。


核心机制与源码级原理

1 exchange() 内部逻辑(简化版)

// JDK 内部实现(深度简化)
public V exchange(V x) throws InterruptedException {
    // 1. 记录当前线程和携带的数据
    Node node = new Node(Thread.currentThread(), x);
    // 2. 尝试把 node 放入槽位(CAS)
    if (casSlot(node)) {
        // 3. 等待另一个线程到达(park)
        LockSupport.park();
        // 4. 唤醒后,从配对节点取出对方的数据
        return node.match;
    } else {
        // 5. 已经有线程在等待,则直接交换数据并唤醒对方
        Node other = slot.get();
        other.match = x;
        LockSupport.unpark(other.waiter);
        return other.item;
    }
}

关键点:线程 A 先到达则自旋/阻塞等待线程 B;线程 B 到达后直接读取 A 的数据并唤醒 A,这保证了数据交换的原子性和双向性。

2 与 CountDownLatch/CyclicBarrier 的区别

工具 线程数 数据流向 典型用
Exchanger 固定2个 双向交换 双缓冲、配对
CyclicBarrier N个 无数据流动 多线程等待齐步走
CountDownLatch 1个等待N个 单向通知 等待初始化完成

实战案例 1:经典生产者-消费者双向握手

场景描述

有一个数据处理流水线:线程A 负责从数据库读取原始订单,线程B 负责对订单进行富化(补充用户信息),两个线程需要不断交换数据:A 把“原始订单”给 B,B 把“富化后订单”返回给 A 用于最终入库。

完整代码实现

import java.util.concurrent.Exchanger;
public class OrderExchangeDemo {
    static class Order {
        String id;
        String status;
        Order(String id, String status) { this.id = id; this.status = status; }
        @Override public String toString() { return "Order[" + id + ":" + status + "]"; }
    }
    public static void main(String[] args) {
        Exchanger<Order> exchanger = new Exchanger<>();
        // 线程A:读取原始订单
        Thread reader = new Thread(() -> {
            for (int i = 1; i <= 3; i++) {
                try {
                    Order raw = new Order("RAW" + i, "PENDING");
                    System.out.println("[A] 读取原始订单: " + raw);
                    // 把原始订单交给B,等待B返回富化订单
                    Order enriched = exchanger.exchange(raw);
                    System.out.println("[A] 收到富化订单: " + enriched);
                } catch (InterruptedException e) {
                    Thread.currentThread().interrupt();
                }
            }
        }, "Reader");
        // 线程B:富化订单
        Thread enricher = new Thread(() -> {
            for (int i = 1; i <= 3; i++) {
                try {
                    Order raw = exchanger.exchange(null); // B 先等待A提供数据
                    if (raw != null) {
                        System.out.println("[B] 收到原始订单: " + raw);
                        Order enriched = new Order(raw.id, "ENRICHED_" + raw.id);
                        // 交换富化结果给A
                        exchanger.exchange(enriched);
                    }
                } catch (InterruptedException e) {
                    Thread.currentThread().interrupt();
                }
            }
        }, "Enricher");
        reader.start();
        enricher.start();
    }
}

运行输出(部分)

[A] 读取原始订单: Order[RAW1:PENDING]
[A] 收到富化订单: Order[RAW1:ENRICHED_RAW1]
...

注意exchange 必须成对出现,B 线程先执行 exchange(null),A 线程后执行 exchange(raw),则 B 会拿到 A 的 raw 对象,A 会拿到 B 的 null,本例通过交替交换(A先给,B后给)实现了方向控制。


实战案例 2:双缓冲(Double Buffering)并发优化

场景描述

一个高频日志系统,写入线程不断产生日志,输出线程负责将日志批量写入磁盘,如果直接用 ArrayList 加锁,会导致大量竞争,使用 Exchanger 可实现双缓冲:写入线程填满缓冲区A后,和输出线程交换——写入线程继续填缓冲区B,输出线程处理缓冲区A。

代码实现

import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.Exchanger;
public class DoubleBufferDemo {
    private static final int BUFFER_SIZE = 10000;
    public static void main(String[] args) {
        Exchanger<List<String>> exchanger = new Exchanger<>();
        List<String> writeBuffer = new ArrayList<>(BUFFER_SIZE);
        List<String> readBuffer = new ArrayList<>(BUFFER_SIZE);
        // 写入线程
        Thread writer = new Thread(() -> {
            int count = 0;
            while (true) {
                writeBuffer.add("LogEntry-" + count++);
                if (writeBuffer.size() >= BUFFER_SIZE) {
                    try {
                        // 把满的写缓冲交给读取线程,换取空缓冲
                        writeBuffer = exchanger.exchange(writeBuffer);
                        writeBuffer.clear(); // 清空新拿到的空缓冲
                    } catch (InterruptedException e) { Thread.currentThread().interrupt(); }
                }
            }
        });
        // 输出线程
        Thread output = new Thread(() -> {
            while (true) {
                try {
                    // 换取写线程的满缓冲
                    readBuffer = exchanger.exchange(readBuffer);
                    System.out.println("输出线程处理 " + readBuffer.size() + " 条日志");
                    // 模拟写入磁盘
                    readBuffer.clear();
                } catch (InterruptedException e) { Thread.currentThread().interrupt(); }
            }
        });
        writer.start();
        output.start();
    }
}

性能优势

  • 零锁:交换操作基于 CAS,不需要 synchronized
  • 无锁内存可见性exchange 内部有 volatilefinal 保证,交换后数据对其他线程立即可见。
  • 吞吐量:相比 LinkedBlockingQueue,双缓冲减少内存分配和GC压力。

实战案例 3:遗传算法中的种群交叉配对

场景描述

在遗传算法中,两个种群(种群A和种群B)需要每次迭代后交换部分优秀个体,使用 Exchanger 可以让两个种群的计算线程在迭代边界自动交换基因序列。

代码片段

import java.util.Random;
import java.util.concurrent.Exchanger;
public class GeneticAlgorithmDemo {
    static class Population {
        int[] genes;
        Population(int size) { genes = new int[size]; }
    }
    public static void main(String[] args) {
        Exchanger<Population> exchanger = new Exchanger<>();
        int populationSize = 50;
        // 种群A 计算线程
        Thread popA = new Thread(() -> {
            Population p = new Population(populationSize);
            for (int gen = 0; gen < 100; gen++) {
                // 模拟适应度计算和选择
                for (int i = 0; i < p.genes.length; i++) {
                    p.genes[i] = new Random().nextInt(100);
                }
                try {
                    if (gen % 10 == 0) { // 每10代交换一次
                        p = exchanger.exchange(p);
                        System.out.println("种群A 交换后: " + p.genes[0] + " ...");
                    }
                } catch (InterruptedException e) { Thread.currentThread().interrupt(); }
            }
        });
        Thread popB = new Thread(() -> {
            Population p = new Population(populationSize);
            for (int gen = 0; gen < 100; gen++) {
                // 种群B 的进化逻辑
                try {
                    if (gen % 10 == 0) {
                        p = exchanger.exchange(p);
                    }
                } catch (InterruptedException e) { Thread.currentThread().interrupt(); }
            }
        });
        popA.start();
        popB.start();
    }
}

注意:此例强调 Exchanger 在非对称业务逻辑中的灵活性——两个线程不必每次迭代都交换,可以按业务条件控制。


高频面试问答与踩坑指南

Q1:Exchanger 中如果有一个线程永远不调用 exchange(),会发生什么?

回答:另一个线程会永久阻塞(除非使用超时重载),如果不希望无限等待,必须使用 exchange(x, 1, TimeUnit.SECONDS),底层通过 LockSupport.parkNanos 实现,超时后抛出 TimeoutException

Q2:Exchanger 能支持两个以上线程吗?

回答:不能,Exchanger 设计上只支持两个线程交换,如果需要多线程(N个)两两互相交换数据,可以使用 Exchanger 配合 Exchanger.Group(JDK 9+),或者改用 ConcurrentHashMap 手动管理配对,但在 JDK 8 及以下,多线程场景建议使用 CyclicBarrier+共享容器。

Q3:exchange() 方法抛出 InterruptedException 时,应该如何处理?

回答:建议 恢复中断状态():

try {
    exchanger.exchange(data);
} catch (InterruptedException e) {
    Thread.currentThread().interrupt(); // 重设中断标志
    // 根据业务决定退出或重试
}

注意:不要吞掉中断异常,否则外层代码无法感知线程中断状态。

Q4:Exchanger 交换的是引用还是值?

回答:交换的是 对象引用,不会深拷贝对象,因此交换后双方持有同一个对象,修改会影响对方,如果不想共享可变状态,需要自行拷贝。

Q5:生产环境中,Exchanger 的性能一定比 BlockingQueue 好吗?

回答不一定,在少量数据、固定双线程场景下,Exchanger 更快(无锁),但:

  • 如果线程数量不固定,Exchanger 会导致部分线程无限等待。
  • 如果交换频率低,CAS 自旋开销可能高于队列。
  • 官方 Javadoc 也指出:Exchanger 适合缓冲区大小动态或无法预知的场景。

总结与最佳实践

  • 使用 Exchanger 的 3 个硬性条件

    1. 线程数量永远且只有 2 个
    2. 交换动作必须成对出现(不能一边有 while,一边只有 if)。
    3. 对数据实时性要求高,且愿意承担阻塞风险。
  • 推荐组合Exchanger + 超时控制 + 循环重试,防止死锁。

  • 与 CompletableFuture 的对比:如果只是单向数据传递,用 CompletableFuture 更简单;只有双向同步交换才用 Exchanger。

请记住 Exchanger 的最经典一句话总结:它让两个线程在约定的汇合点,互相把对方需要的东西交给对方,然后继续前行——这是并发设计中最优雅的握手之一。

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