Java生产者消费者模式实战:从BlockingQueue到虚拟线程的完整案例解析
📚 目录导读
- 为什么生产者消费者模式是并发编程的基石?
- 经典实现一:
wait()/notifyAll()手写同步协作 - 经典实现二:
BlockingQueue一行代码解决核心问题 - 进阶挑战:多生产者/多消费者下的负载均衡
- 性能优化:使用虚拟线程(Project Loom)重写案例
- 高频面试问答:深度解析
sleep与wait、死锁规避等 - 总结与最佳实践建议
为什么生产者消费者模式是并发编程的基石?
在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:BlockingQueue 的 put 和 offer 有何区别?
put阻塞等待,直到成功放入队列。offer尝试放入,若队列满则立即返回false(或等待指定时间)。
总结与最佳实践建议
- 单生产者/消费者下,使用
wait/notify可以加深对JMM内存模型的理解。 - 生产环境中,优先使用
BlockingQueue,它更安全、更高效。 - 考虑数据丢失:若消费者处理失败,需要设计重试机制或死信队列。
- 监控:定期打印队列大小和线程状态,排查性能瓶颈。
最后给出一个完整可运行的示例(基于BlockingQueue)供读者二次封装,并建议结合JUnit进行并发压力测试,确保无数据丢失、无重复消费。并发编程的黄金法则是:能不用锁就不用锁,必须用锁时尽量缩小锁范围,希望此文能助你在代码实践中游刃有余。