Java生产者消费者案例代码

wen java案例 3

Java生产者消费者模式实战:从BlockingQueue到虚拟线程的完整案例解析


📚 目录导读

  1. 为什么生产者消费者模式是并发编程的基石?
  2. 经典实现一:wait()/notifyAll() 手写同步协作
  3. 经典实现二:BlockingQueue 一行代码解决核心问题
  4. 进阶挑战:多生产者/多消费者下的负载均衡
  5. 性能优化:使用虚拟线程(Project Loom)重写案例
  6. 高频面试问答:深度解析sleepwait、死锁规避等
  7. 总结与最佳实践建议

为什么生产者消费者模式是并发编程的基石?

在Java并发编程中,生产者-消费者模式 是解耦数据生产与消费逻辑的核心架构,它解决了两个关键问题:

Java生产者消费者案例代码

  • 速度不匹配:生产者生成数据的速度可能远超消费者的处理能力(或反之),导致资源浪费或数据丢失。
  • 时序依赖:生产者无需等待消费者处理完上一个数据才能继续生产,两者通过缓冲区(如队列)进行异步通信。

该模式不仅用于线程池任务队列,还广泛应用于消息中间件(如Kafka)、日志采集系统等,掌握其代码实现,是深入理解JUC(Java并发包)的必经之路。


经典实现一:wait()/notifyAll() 手写同步协作

核心逻辑:使用一个共享的LinkedList作为缓冲区,通过synchronized锁保证线程安全,当缓冲区满时,生产者线程调用wait()进入等待;当缓冲区空时,消费者线程调用wait(),每次操作后调用notifyAll()唤醒对方。

public class ProducerConsumerWaitNotify {
    private static final int CAPACITY = 5;
    private final LinkedList<Integer> queue = new LinkedList<>();
    public synchronized void produce(int value) throws InterruptedException {
        while (queue.size() == CAPACITY) {
            wait(); // 缓冲区满,等待消费者消费
        }
        queue.add(value);
        System.out.println("生产:" + value + ",当前大小:" + queue.size());
        notifyAll(); // 唤醒可能等待的消费者
    }
    public synchronized int consume() throws InterruptedException {
        while (queue.isEmpty()) {
            wait(); // 缓冲区空,等待生产者生产
        }
        int value = queue.removeFirst();
        System.out.println("消费:" + value + ",剩余大小:" + queue.size());
        notifyAll();
        return value;
    }
    // 测试主方法(略)
}

注意:必须使用while而非if进行条件判断,防止虚假唤醒(spurious wakeup)导致越界错误。


经典实现二:BlockingQueue 一行代码解决核心问题

java.util.concurrent.BlockingQueue 接口提供了内置的阻塞方法(put()take()),它们自动处理锁和等待通知机制,极大地简化了代码

public class ProducerConsumerBlockingQueue {
    private static final BlockingQueue<Integer> queue = new LinkedBlockingQueue<>(5);
    static class Producer implements Runnable {
        public void run() {
            try {
                int i = 0;
                while (true) {
                    queue.put(i++); // 自动阻塞直到有空间
                    Thread.sleep(100);
                }
            } catch (InterruptedException e) { Thread.currentThread().interrupt(); }
        }
    }
    static class Consumer implements Runnable {
        public void run() {
            try {
                while (true) {
                    Integer data = queue.take(); // 自动阻塞直到有元素
                    System.out.println("消费:" + data);
                }
            } catch (InterruptedException e) { Thread.currentThread().interrupt(); }
        }
    }
    // 启动线程代码(略)
}

优势LinkedBlockingQueue 内部使用两把锁(takeLock和putLock),提高了吞吐量,且无需手工处理wait/notify,代码更加健壮。


进阶挑战:多生产者/多消费者下的负载均衡

当有多个生产者(线程A、B)和多个消费者(线程C、D)时,单纯的notifyAll()可能导致惊群效应(所有线程同时唤醒,但只有一个能执行),此时建议:

  • 使用ReentrantLock + 多个Condition(生产者条件/消费者条件)实现精准唤醒。
  • 或者直接使用BlockingQueue,它天生支持多线程并发安全。

示例:使用ExecutorService创建固定线程池,提交多个生产者/消费者任务,缓冲区容量不宜过小,防止频繁阻塞;也不宜过大,防止内存溢出。


性能优化:使用虚拟线程(Project Loom)重写案例

Java 21+引入了虚拟线程,它们是轻量级线程,阻塞成本极低,在上述例子中,即使有上千个生产者/消费者,也不会消耗大量OS线程。

// 使用虚拟线程启动生产者(伪代码)
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
    for (int i = 0; i < 10; i++) {
        executor.submit(new Producer()); // 每个生产者是一个虚拟线程
    }
    // 类似地启动消费者
}

优化点:虚拟线程在queue.put()阻塞时,会自动释放载体线程(Carrier Thread),使得系统能轻松支撑高并发IO等待场景,这是传统平台线程无法比拟的。


高频面试问答:深度解析关键细节

问1:sleep()wait()有什么区别?

  • sleep不释放锁,wait释放锁,生产消费场景必须用wait,否则会死锁。
  • sleep是静态方法,wait是Object实例方法。

问2:如何避免死锁?

  • 保证加锁顺序一致(如先锁producerLock再锁consumerLock)。
  • 使用tryLock带超时时间。
  • wait循环中使用超时(wait(1000)),防止永久等待。

问3:BlockingQueueputoffer 有何区别?

  • put 阻塞等待,直到成功放入队列。
  • offer 尝试放入,若队列满则立即返回false(或等待指定时间)。

总结与最佳实践建议

  • 单生产者/消费者下,使用wait/notify可以加深对JMM内存模型的理解。
  • 生产环境中,优先使用BlockingQueue,它更安全、更高效。
  • 考虑数据丢失:若消费者处理失败,需要设计重试机制或死信队列。
  • 监控:定期打印队列大小和线程状态,排查性能瓶颈。

最后给出一个完整可运行的示例(基于BlockingQueue)供读者二次封装,并建议结合JUnit进行并发压力测试,确保无数据丢失、无重复消费。并发编程的黄金法则是:能不用锁就不用锁,必须用锁时尽量缩小锁范围,希望此文能助你在代码实践中游刃有余。

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