跳到主要内容

BlockingQueue 与生产消费

BlockingQueue 在 Queue 基础上增加等待语义:队列为空时消费者可以等待,队列已满时生产者可以等待。使用有界队列还能把下游处理能力转化为明确的背压信号。

1. 四组入队与出队方法

BlockingQueue 为同一操作提供不同失败策略:

操作抛异常返回特殊值一直等待限时等待
插入addofferputoffer(e, timeout, unit)
取出removepolltakepoll(timeout, unit)
查看队首elementpeek不适用不适用
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 不存储元素,只做生产者和消费者的直接交接。

选择要看是否需要缓冲、允许多少积压、对象分配和吞吐要求。服务端通常显式设置有界容量,并为满队列定义超时或拒绝策略。