本文目录导读:

这是一个经典的Java并发编程案例,阻塞队列(BlockingQueue)是一种支持生产者-消费者模式的线程安全队列,当队列为空时,消费者线程会阻塞等待;当队列满时,生产者线程会阻塞等待。
下面我将从实现原理、手写实现(基于ReentrantLock和Condition)以及使用Java内置API三个维度来讲解。
实现原理
阻塞队列的核心需要解决两个问题:
- 线程安全:多个线程同时读写队列时数据不混乱(需要锁保护)
- 等待/通知机制:队列满时生产者等待,队列空时消费者等待(需要条件变量)
Java中实现上述机制主要有两种方式:
- synchronized + wait/notify
- ReentrantLock + Condition(更灵活,推荐)
下面我们手写一个基于ReentrantLock + Condition的有界阻塞队列。
手写实现(面试常考)
import java.util.LinkedList;
import java.util.concurrent.locks.Condition;
import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock;
public class MyBlockingQueue<T> {
// 底层存储结构
private final LinkedList<T> queue;
// 队列容量
private final int capacity;
// 锁
private final Lock lock;
// 队列非空条件(消费者等待)
private final Condition notEmpty;
// 队列非满条件(生产者等待)
private final Condition notFull;
public MyBlockingQueue(int capacity) {
this.capacity = capacity;
this.queue = new LinkedList<>();
this.lock = new ReentrantLock();
this.notEmpty = lock.newCondition();
this.notFull = lock.newCondition();
}
/**
* 往队列尾部添加元素,如果队列满则阻塞等待
*/
public void put(T element) throws InterruptedException {
lock.lock(); // 获取锁(可中断)
try {
// 当队列满时,生产者线程在notFull条件上等待
while (queue.size() == capacity) {
notFull.await();
}
// 生产数据
queue.addLast(element);
// 通知消费者:队列非空了
notEmpty.signal();
} finally {
lock.unlock();
}
}
/**
* 从队列头部取出元素,如果队列空则阻塞等待
*/
public T take() throws InterruptedException {
lock.lock();
try {
// 当队列空时,消费者线程在notEmpty条件上等待
while (queue.size() == 0) {
notEmpty.await();
}
// 消费数据
T element = queue.removeFirst();
// 通知生产者:队列非满了
notFull.signal();
return element;
} finally {
lock.unlock();
}
}
/**
* 返回队列当前大小(调试用)
*/
public int size() {
lock.lock();
try {
return queue.size();
} finally {
lock.unlock();
}
}
}
使用示例(生产者-消费者)
public class BlockingQueueDemo {
public static void main(String[] args) {
// 创建一个容量为3的阻塞队列
MyBlockingQueue<Integer> queue = new MyBlockingQueue<>(3);
// 生产者线程
Thread producer = new Thread(() -> {
try {
for (int i = 1; i <= 10; i++) {
queue.put(i);
System.out.println("生产: " + i);
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
});
// 消费者线程
Thread consumer = new Thread(() -> {
try {
for (int i = 1; i <= 10; i++) {
Integer value = queue.take();
System.out.println("消费: " + value);
Thread.sleep(100); // 模拟消费耗时
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
});
producer.start();
consumer.start();
}
}
关键设计要点
为什么用 while 而不是 if 进行条件判断?
// 必须使用 while 循环,不能使用 if
while (queue.size() == capacity) {
notFull.await();
}
原因:线程被唤醒后,条件可能再次不满足(虚假唤醒或另一个生产者抢先填充了队列),使用 while 循环可以重新检查条件,这是 Guarded Suspension(保护性暂停) 模式的要求。
signal() vs signalAll()
signal():只唤醒一个等待线程(更高效),当前场景中,生产者只唤醒一个消费者,消费者只唤醒一个生产者,所以用signal()足够了。signalAll():唤醒所有等待线程(安全性更高但性能较差)。
锁的选择
| 特性 | synchronized | ReentrantLock |
|---|---|---|
| 可中断 | 不支持 | 支持 |
| 超时等待 | 不支持 | 支持 |
| 多个条件变量 | 不直接支持 | 支持(newCondition()) |
| 公平性 | 非公平 | 可设置公平/非公平 |
对于阻塞队列这种需要两个条件变量(notEmpty和notFull)的场景,ReentrantLock + Condition 是更自然的选择。
使用Java内置的BlockingQueue
JDK已经提供了完善的阻塞队列实现,实际开发中直接使用即可:
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.BlockingQueue;
public class BuiltInBlockingQueueDemo {
public static void main(String[] args) throws InterruptedException {
// 创建容量为3的有界阻塞队列
BlockingQueue<Integer> queue = new ArrayBlockingQueue<>(3);
// 生产者
new Thread(() -> {
try {
for (int i = 1; i <= 10; i++) {
queue.put(i);
System.out.println("生产: " + i);
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}).start();
// 消费者
new Thread(() -> {
try {
for (int i = 1; i <= 10; i++) {
Integer value = queue.take();
System.out.println("消费: " + value);
Thread.sleep(100);
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}).start();
}
}
Java内置的阻塞队列实现(基于 ReentrantLock):
| 实现类 | 特性 |
|---|---|
ArrayBlockingQueue |
有界、数组结构、FIFO |
LinkedBlockingQueue |
可有界/无界、链表结构 |
PriorityBlockingQueue |
无界、支持优先级排序 |
SynchronousQueue |
容量为0(直接传递) |
面试扩展:用synchronized实现
如果面试官不让用ReentrantLock,可以用 synchronized + wait/notify:
public class SyncBlockingQueue<T> {
private final LinkedList<T> queue = new LinkedList<>();
private final int capacity;
public SyncBlockingQueue(int capacity) {
this.capacity = capacity;
}
public synchronized void put(T element) throws InterruptedException {
while (queue.size() == capacity) {
wait();
}
queue.addLast(element);
notifyAll(); // 必须用 notifyAll,防止只唤醒同类线程
}
public synchronized T take() throws InterruptedException {
while (queue.isEmpty()) {
wait();
}
T element = queue.removeFirst();
notifyAll();
return element;
}
}
注意:用 synchronized 时,只有一个条件变量(wait/notify 是全局的),所以必须使用 notifyAll(),否则可能发生“信号丢失”(生产者唤醒了生产者,消费者没被唤醒)。
| 维度 | 说明 |
|---|---|
| 核心机制 | 锁 + 条件等待/通知 |
| 推荐实现 | ReentrantLock + 两个Condition |
| 关键细节 | while 循环检查条件、signal() 精确唤醒 |
| 生产使用 | 直接使用 ArrayBlockingQueue / LinkedBlockingQueue |
如果这是面试场景,建议按以下思路回答:
- 先讲清楚阻塞队列核心功能(满则等、空则等)
- 手写核心代码(put/take)
- 解释为什么用while不用if
- 对比synchronized和ReentrantLock方案的区别