BlockingQueue 全过程原理
BlockingQueue 是 Java 并发里连接“生产者”和“消费者”的核心工具,也是线程池任务队列、异步削峰、批处理缓冲、定时延迟任务的基础。它不只是一个队列,更重要的是:队列满了生产者怎么等,队列空了消费者怎么等,系统压力来了怎么背压。
学习目标
| 目标 | 要能说清楚 |
|---|---|
| 基本模型 | 生产者、队列、消费者如何解耦 |
| 阻塞原理 | put/take 为什么会等待,谁来唤醒 |
| API 区别 | add/offer/put/take/poll 的行为差异 |
| 队列选型 | ArrayBlockingQueue、LinkedBlockingQueue、SynchronousQueue、PriorityBlockingQueue、DelayQueue 怎么选 |
| 线程池关系 | 不同队列如何影响线程池扩容和拒绝策略 |
| 商业场景 | 异步入库、通知发送、采集缓冲、延迟取消 |
| 生产排查 | 队列堆积、消费者退出、无界队列 OOM、忽略中断怎么查 |
为什么需要 BlockingQueue
没有队列时,生产者必须直接调用消费者:
flowchart TD
A["请求线程"] --> B["直接调用短信服务"]
B --> C["直接调用邮件服务"]
C --> D["直接写数据库"]问题:
| 问题 | 后果 |
|---|---|
| 下游慢 | 请求线程被拖慢 |
| 瞬时流量高 | 下游被打爆 |
| 消费失败 | 请求链路复杂,重试难做 |
| 生产和消费强耦合 | 一个环节慢影响整个链路 |
用 BlockingQueue 后:
flowchart TD
A["生产者"] --> B["BlockingQueue"]
B --> C["消费者1"]
B --> D["消费者2"]
B --> E["消费者3"]
B --> F["队列满时生产者等待或失败"]
B --> G["队列空时消费者等待"]它带来的核心价值:
| 价值 | 说明 |
|---|---|
| 解耦 | 生产者只负责放任务,消费者异步处理 |
| 削峰 | 短时间高峰先进入队列 |
| 背压 | 队列满时让生产者等待、失败或降级 |
| 平滑处理 | 消费者按稳定速度处理任务 |
API 行为差异
BlockingQueue 的方法很多,不同方法在队列满/空时行为不同。
| 操作 | 抛异常 | 返回特殊值 | 一直阻塞 | 超时等待 |
|---|---|---|---|---|
| 插入 | add(e) | offer(e) | put(e) | offer(e, time, unit) |
| 移除 | remove() | poll() | take() | poll(time, unit) |
| 查看 | element() | peek() | 无 | 无 |
怎么选:
| 场景 | 推荐 |
|---|---|
| 必须提交,队列满就等 | put |
| 不能无限等,需要快速失败 | offer |
| 可以等一会,超时降级 | offer(timeout) |
| 消费者一直等任务 | take |
| 消费者定时检查退出条件 | poll(timeout) |
put 阻塞流程
以有界队列为例:
flowchart TD
A["生产者调用 put"] --> B["获取入队锁"]
B --> C{"队列是否已满"}
C -->|否| D["元素入队"]
D --> E["唤醒等待取数据的消费者"]
E --> F["释放锁"]
C -->|是| G["生产者进入 notFull 条件队列等待"]
G --> H["消费者取走元素后 signal notFull"]
H --> B核心点:
| 点 | 说明 |
|---|---|
| 队列满 | 生产者不能继续放,否则内存无限增长 |
| notFull | 表示“队列未满”的等待条件 |
| signal | 消费者取走元素后,唤醒一个等待入队的生产者 |
| 中断 | 阻塞等待时可以响应 interrupt |
take 阻塞流程
flowchart TD
A["消费者调用 take"] --> B["获取出队锁"]
B --> C{"队列是否为空"}
C -->|否| D["取出元素"]
D --> E["唤醒等待放数据的生产者"]
E --> F["释放锁"]
C -->|是| G["消费者进入 notEmpty 条件队列等待"]
G --> H["生产者放入元素后 signal notEmpty"]
H --> B这就是 BlockingQueue 名字里“Blocking”的含义:条件不满足时线程进入等待,不是一直死循环占 CPU。
ArrayBlockingQueue
ArrayBlockingQueue 是基于数组的有界阻塞队列。
flowchart TD
A["固定容量数组"] --> B["putIndex 入队位置"]
A --> C["takeIndex 出队位置"]
B --> D["循环递增"]
C --> D特点:
| 特点 | 说明 |
|---|---|
| 有界 | 创建时必须指定容量 |
| 数组结构 | 内存更连续 |
| 单锁实现 | 入队和出队共用一把锁 |
| 可选公平锁 | 构造参数可控制公平性 |
Demo:
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.BlockingQueue;
public class ArrayBlockingQueueDemo {
public static void main(String[] args) throws Exception {
BlockingQueue<String> queue = new ArrayBlockingQueue<>(2);
queue.put("task-1");
queue.put("task-2");
boolean success = queue.offer("task-3");
System.out.println(success); // false
}
}适合固定容量、需要明确背压的任务队列。
LinkedBlockingQueue
LinkedBlockingQueue 是基于链表的阻塞队列,可以有界,也可以近似无界。
特点:
| 特点 | 说明 |
|---|---|
| 链表结构 | 节点按需创建 |
| 可指定容量 | 不指定时容量非常大 |
| 两把锁 | putLock 和 takeLock 分离,入队出队可并发 |
| 常用于线程池 | newFixedThreadPool 默认使用无界 LinkedBlockingQueue |
危险点:
new LinkedBlockingQueue<>();这不是“安全无限队列”。任务生产速度长期大于消费速度时,队列会持续增长,最终可能 OOM。
生产建议:
new LinkedBlockingQueue<>(1000);容量必须结合内存、任务大小、处理速度和可接受延迟来估算。
SynchronousQueue
SynchronousQueue 不存储元素,它是直接交接队列。
flowchart TD
A["生产者 put"] --> B{"是否有消费者正在 take"}
B -->|有| C["直接交给消费者"]
B -->|没有| D["生产者等待"]特点:
| 特点 | 说明 |
|---|---|
| 容量为 0 | 不保存任务 |
| 直接交接 | 生产者和消费者必须配对 |
| 常用于 cached 线程池 | 没有空闲线程时倾向创建新线程 |
这解释了为什么 newCachedThreadPool 可能创建很多线程:它用 SynchronousQueue,不排队,任务交不出去就尝试创建新线程。
PriorityBlockingQueue
优先级阻塞队列,元素按优先级出队。
import java.util.concurrent.PriorityBlockingQueue;
public class PriorityQueueDemo {
static class Task implements Comparable<Task> {
final int priority;
final String name;
Task(int priority, String name) {
this.priority = priority;
this.name = name;
}
@Override
public int compareTo(Task other) {
return Integer.compare(other.priority, this.priority);
}
}
public static void main(String[] args) throws Exception {
PriorityBlockingQueue<Task> queue = new PriorityBlockingQueue<>();
queue.put(new Task(1, "low"));
queue.put(new Task(10, "high"));
System.out.println(queue.take().name); // high
}
}注意:它默认无界,不能靠它天然背压。优先级任务也可能导致低优先级任务长期饥饿。
DelayQueue
DelayQueue 只有元素到期后才能被取出,常用于超时任务、延迟取消。
import java.util.concurrent.DelayQueue;
import java.util.concurrent.Delayed;
import java.util.concurrent.TimeUnit;
public class DelayQueueDemo {
static class DelayTask implements Delayed {
private final String orderId;
private final long executeAt;
DelayTask(String orderId, long delayMs) {
this.orderId = orderId;
this.executeAt = System.currentTimeMillis() + delayMs;
}
@Override
public long getDelay(TimeUnit unit) {
long delay = executeAt - System.currentTimeMillis();
return unit.convert(delay, TimeUnit.MILLISECONDS);
}
@Override
public int compareTo(Delayed other) {
return Long.compare(this.getDelay(TimeUnit.MILLISECONDS),
other.getDelay(TimeUnit.MILLISECONDS));
}
}
public static void main(String[] args) throws Exception {
DelayQueue<DelayTask> queue = new DelayQueue<>();
queue.put(new DelayTask("order-1", 1000));
System.out.println(queue.take().orderId);
}
}商业系统里,订单超时取消更常用 MQ 延迟消息或定时任务扫描,因为 DelayQueue 是 JVM 内存级别,服务重启会丢任务,分布式部署也难统一。
线程池和队列的关系
线程池提交任务的关键流程:
flowchart TD
A["提交任务"] --> B{"工作线程数 < corePoolSize"}
B -->|是| C["创建核心线程执行"]
B -->|否| D{"队列是否能入队"}
D -->|能| E["任务进入 BlockingQueue"]
D -->|不能| F{"工作线程数 < maximumPoolSize"}
F -->|是| G["创建非核心线程执行"]
F -->|否| H["执行拒绝策略"]队列会直接影响线程池行为:
| 队列 | 对线程池的影响 |
|---|---|
| 无界 LinkedBlockingQueue | 核心线程满后一直排队,通常不会扩到 maximumPoolSize,风险是 OOM |
| 有界 ArrayBlockingQueue | 队列满后才扩到 maximumPoolSize,最终触发拒绝策略 |
| SynchronousQueue | 不排队,交不出去就创建线程,风险是线程数暴涨 |
| PriorityBlockingQueue | 按优先级执行,但通常无界,且最大线程数可能不按预期发挥 |
所以线程池参数不能只看 core/max,还必须看队列类型和容量。
商业场景:异步通知发送
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.BlockingQueue;
public class NotifyService {
private final BlockingQueue<String> queue = new ArrayBlockingQueue<>(1000);
public boolean submit(String message) {
return queue.offer(message);
}
public void startConsumer() {
Thread consumer = new Thread(() -> {
while (!Thread.currentThread().isInterrupted()) {
try {
String message = queue.take();
send(message);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
} catch (Exception e) {
// 记录失败并进入重试或死信逻辑
e.printStackTrace();
}
}
}, "notify-consumer");
consumer.start();
}
private void send(String message) {
System.out.println("send: " + message);
}
}为什么 submit 用 offer 而不是 put?
接口线程通常不能无限阻塞。队列满时返回 false,调用方可以降级、返回繁忙、记录失败或走 MQ。
商业场景:批量入库缓冲
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.TimeUnit;
public class BatchWriter {
private final BlockingQueue<String> queue = new LinkedBlockingQueue<>(5000);
public boolean submit(String row) {
return queue.offer(row);
}
public void consumeLoop() throws InterruptedException {
List<String> batch = new ArrayList<>(100);
while (!Thread.currentThread().isInterrupted()) {
String first = queue.poll(1, TimeUnit.SECONDS);
if (first == null) {
continue;
}
batch.add(first);
queue.drainTo(batch, 99);
writeToDb(batch);
batch.clear();
}
}
private void writeToDb(List<String> rows) {
System.out.println("batch insert: " + rows.size());
}
}这个模式适合采集数据、日志、指标等批量写入。但要注意:JVM 队列不是可靠消息队列,服务宕机会丢内存中的数据。强可靠场景应使用 MQ 或落库状态表。
线上排查:队列堆积
flowchart TD
A["队列长度上涨"] --> B["比较生产 TPS 和消费 TPS"]
B --> C{"生产是否长期大于消费"}
C -->|是| D["消费者处理能力不足"]
D --> E["查慢 SQL、远程调用、锁、批量大小"]
C -->|否| F["检查消费者是否异常退出"]
F --> G["检查是否有拒绝、重试风暴、下游限流"]排查指标:
| 指标 | 含义 |
|---|---|
| queue size | 当前堆积量 |
| offer fail count | 入队失败次数 |
| consume TPS | 消费速度 |
| task cost P95/P99 | 单任务耗时 |
| consumer alive | 消费线程是否还活着 |
| downstream latency | 下游接口或数据库是否变慢 |
不要一看到堆积就盲目加消费者。如果下游数据库已经慢了,加消费者可能让数据库更慢,堆积短暂下降后又继续上涨。
常见坑
| 坑 | 后果 | 正确做法 |
|---|---|---|
| 使用无界队列 | 堆积到 OOM | 使用有界队列 |
队列满还 put | 请求线程大量阻塞 | 用 offer(timeout) 或快速失败 |
| 消费者异常退出 | 队列只进不出 | 捕获异常并监控线程存活 |
| 忽略 InterruptedException | 线程无法优雅停止 | 恢复中断标记并退出 |
| 单机队列当可靠 MQ | 重启丢任务 | 强可靠用 MQ 或持久化 |
| 盲目加消费者 | 下游被打爆 | 结合下游容量扩容 |
面试标准回答
BlockingQueue 是什么
BlockingQueue 是支持阻塞插入和阻塞获取的线程安全队列,常用于生产者消费者模型。队列满时,put 会等待空间;队列空时,take 会等待元素。它可以解耦生产和消费、削峰、背压,也是线程池任务队列的重要基础。
put 和 offer 区别
put 在队列满时会一直阻塞,适合必须提交且允许等待的场景;offer 在队列满时立即返回 false,适合接口线程快速失败;offer(timeout) 可以等待一段时间,超时后降级。
ArrayBlockingQueue 和 LinkedBlockingQueue 区别
ArrayBlockingQueue 基于数组,必须指定固定容量,入队出队共用一把锁,适合明确有界背压。LinkedBlockingQueue 基于链表,可指定容量,不指定时容量非常大,通常有 putLock 和 takeLock 两把锁,吞吐性能更好,但无界使用有 OOM 风险。
SynchronousQueue 为什么容易让线程数变多
SynchronousQueue 不存储元素,生产者必须直接把任务交给消费者。线程池使用它时,如果没有空闲线程接任务,就倾向于创建新线程,直到 maximumPoolSize。因此 cached 线程池在高流量下可能线程数暴涨。
队列堆积怎么排查
先比较生产 TPS 和消费 TPS,如果生产长期大于消费,队列一定上涨。再查消费者是否存活、单任务耗时、慢 SQL、远程调用、锁等待、下游限流和重试风暴。处理时不能只加消费者,要结合下游容量、队列容量、拒绝策略和降级策略。
关联知识点
| 知识点 | 为什么要看 |
|---|---|
| 并发集合与阻塞队列 | 并发容器整体导读 |
| 线程池执行流程 | 队列如何影响线程池扩容和拒绝 |
| 线程池生命周期与排查 | 队列堆积在线程池中的完整排查 |
| ConcurrentHashMap全过程 | 对比并发 Map 和阻塞队列解决的问题不同 |
| 消息队列总览 | 强可靠异步场景为什么需要 MQ |
