BlockingQueue 与生产消费
BlockingQueue 在 Queue 基础上增加等待语义:队列为空时消费者可以等待,队列已满时生产者可以等待。使用有界队列还能把下游处理能力转化为明确的背压信号。
1. 四组入队与出队方法
BlockingQueue 为同一操作提供不同失败策略:
| 操作 | 抛异常 | 返回特殊值 | 一直等待 | 限时等待 |
|---|---|---|---|---|
| 插入 | add | offer | put | offer(e, timeout, unit) |
| 取出 | remove | poll | take | poll(timeout, unit) |
| 查看队首 | element | peek | 不适用 | 不适用 |
BlockingQueue<Job> queue = new ArrayBlockingQueue<>(1_000);
queue.put(job); // 满时等待,可被中断
Job next = queue.take(); // 空时等待,可被中断
选择方法就是选择失败协议:等待多久、是否允许调用线程被拖慢、超时后由谁处理任务。
2. 有界队列建立容量边界
生产速度长期高于消费速度时,任何缓冲区都会继续增长或拒绝数据。固定容量让这个边界提前显现:
boolean accepted = queue.offer(job, 100, TimeUnit.MILLISECONDS);
if (!accepted) {
rejectOrPersist(job);
}
满队列时可以:
- 阻塞生产者,让上游自然降速。
- 等待有限时间后失败。
- 拒绝低优先级任务。
- 持久化到可靠队列。
- 由调用线程处理,传播压力。
没有统一最佳策略。关键是不能把“暂时放进内存”当作处理能力提升。
3. 生产者与消费者不再手写 wait/notify
final class Worker implements Runnable {
private final BlockingQueue<Job> queue;
Worker(BlockingQueue<Job> queue) {
this.queue = queue;
}
@Override
public void run() {
try {
while (!Thread.currentThread().isInterrupted()) {
Job job = queue.take();
process(job);
}
} catch (InterruptedException exception) {
Thread.currentThread().interrupt();
}
}
}
队列封装了条件等待、唤醒和可见性。一个线程在入队前的动作 happens-before 另一个线程成功取出该元素后的动作,因此任务对象能按队列协议安全传递。
业务仍需处理:
process失败是否重试。- 重试是否保持顺序。
- 任务是否可以重复执行。
- 进程崩溃时内存任务是否允许丢失。
- 停机时等待、排空还是转存。
4. 常见实现怎样选择
4.1 ArrayBlockingQueue
固定大小循环数组,在构造时确定容量。节点分配少,容量边界明确,可选公平锁。生产和消费围绕同一主锁协调。
4.2 LinkedBlockingQueue
链式节点,可指定容量;不指定时上限很大,容易被误当成真正安全的无界队列。当前实现分别使用 put 和 take 锁,生产消费可以有更多并行,但每个元素需要节点分配。
4.3 SynchronousQueue
容量为零。每次 put 必须与一个 take 直接交接,适合把任务立即移交给空闲消费者,不承担缓冲。
4.4 PriorityBlockingQueue
按优先级取出,逻辑上无界。它不会因为“队列满”给生产者背压;相同优先级也不自动保证 FIFO,任务比较器必须稳定。
4.5 DelayQueue
只有延迟到期的元素才能取出,适合进程内延迟触发。它不持久化,系统时钟、进程重启和大量取消任务都需要考虑。
5. 中断、关闭与 poison pill
BlockingQueue 没有统一 close 方法。常见停止方式:
- 中断消费者线程。
- 发送特殊结束元素,即 poison pill。
- 外部关闭标志配合限时 poll。
- 由拥有线程的 ExecutorService 执行关闭协议。
poison pill 需要保证:
- 结束值不会与普通业务任务混淆。
- 多个消费者收到足够数量的结束信号。
- 队列已满时仍能完成停机。
- 不会越过仍需处理的任务。
线程池任务队列通常由 ExecutorService 管理,不应绕过线程池直接拿内部队列实现自定义关闭。
6. 内存队列与消息队列的边界
BlockingQueue 适合同一 JVM 内线程协作。它不提供:
- 进程崩溃后的持久化恢复。
- 跨进程消费。
- 消费确认、死信和重放。
- 分区扩展与副本容错。
任务不能丢、需要跨服务或要保留较长时间时,应使用持久消息系统或数据库任务表。进程内 BlockingQueue 仍可作为消费者内部的小型有界缓冲,但要与外部确认协议协调。
7. 常见问题
7.1 LinkedBlockingQueue 不传容量有什么风险
它的容量上限非常大。消费者变慢时,生产者不会及时获得压力信号,节点会持续占用堆,最终可能触发长 GC 或 OOM。服务端线程池和任务管道通常显式设置容量。
7.2 remainingCapacity() 后再 put 能保证不阻塞吗
不能。两个调用之间其他线程可能填满队列。直接使用 offer 或限时 offer,把检查与插入交给队列的原子操作。
7.3 BlockingQueue 保证业务只处理一次吗
不保证。消费者可能在业务提交后、记录完成前失败,调用方也可能重试入队。需要幂等键、状态记录或事务消息等业务协议。
7.4 公平队列一定更好吗
公平模式减少长期等待偏差,但通常降低吞吐。只有等待公平是明确需求时开启,并通过负载测试验证。
8. 面试题
8.1 怎样用 BlockingQueue 实现生产者消费者
出现公司:快手
考察重点
- put/take 的等待与中断语义。
- 有界容量怎样形成背压。
- 任务失败、停机与丢失边界。
相关内容:第 1 节“四组入队与出队方法”至第 5 节“中断、关闭与 poison pill”。
参考回答
生产者通过 put 或限时 offer 写入有界 BlockingQueue,消费者循环调用 take 取出并处理。队列满时生产者等待或按超时策略失败,队列空时消费者等待,队列同时建立了任务发布的 happens-before 关系。
还要定义中断和关闭:消费者收到中断后恢复中断标记并退出,或使用明确的结束元素。生产失败、消费重试、幂等和进程崩溃后的任务丢失不由内存队列自动解决。
8.2 ArrayBlockingQueue、LinkedBlockingQueue 和 SynchronousQueue 怎样选择
出现公司:去哪儿、快手
考察重点
- 固定数组、链式节点和零容量交接。
- 容量、分配、锁竞争与背压的关系。
- 默认近似无界为何危险。
相关内容:第 2 节“有界队列建立容量边界”、第 4 节“常见实现怎样选择”。
参考回答
ArrayBlockingQueue 是固定容量循环数组,节点分配少且边界清楚;LinkedBlockingQueue 使用链式节点,可指定容量,生产和消费使用不同锁,但不指定容量时容易积压到内存压力;SynchronousQueue 不存储元素,只做生产者和消费者的直接交接。
选择要看是否需要缓冲、允许多少积压、对象分配和吞吐要求。服务端通常显式设置有界容量,并为满队列定义超时或拒绝策略。