Skip to content

BlockingQueue 全过程原理

BlockingQueue 是 Java 并发里连接“生产者”和“消费者”的核心工具,也是线程池任务队列、异步削峰、批处理缓冲、定时延迟任务的基础。它不只是一个队列,更重要的是:队列满了生产者怎么等,队列空了消费者怎么等,系统压力来了怎么背压。

学习目标

目标要能说清楚
基本模型生产者、队列、消费者如何解耦
阻塞原理put/take 为什么会等待,谁来唤醒
API 区别add/offer/put/take/poll 的行为差异
队列选型ArrayBlockingQueue、LinkedBlockingQueue、SynchronousQueue、PriorityBlockingQueue、DelayQueue 怎么选
线程池关系不同队列如何影响线程池扩容和拒绝策略
商业场景异步入库、通知发送、采集缓冲、延迟取消
生产排查队列堆积、消费者退出、无界队列 OOM、忽略中断怎么查

为什么需要 BlockingQueue

没有队列时,生产者必须直接调用消费者:

mermaid
flowchart TD
    A["请求线程"] --> B["直接调用短信服务"]
    B --> C["直接调用邮件服务"]
    C --> D["直接写数据库"]

问题:

问题后果
下游慢请求线程被拖慢
瞬时流量高下游被打爆
消费失败请求链路复杂,重试难做
生产和消费强耦合一个环节慢影响整个链路

用 BlockingQueue 后:

mermaid
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 阻塞流程

以有界队列为例:

mermaid
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 阻塞流程

mermaid
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 是基于数组的有界阻塞队列。

mermaid
flowchart TD
    A["固定容量数组"] --> B["putIndex 入队位置"]
    A --> C["takeIndex 出队位置"]
    B --> D["循环递增"]
    C --> D

特点:

特点说明
有界创建时必须指定容量
数组结构内存更连续
单锁实现入队和出队共用一把锁
可选公平锁构造参数可控制公平性

Demo:

java
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

危险点:

java
new LinkedBlockingQueue<>();

这不是“安全无限队列”。任务生产速度长期大于消费速度时,队列会持续增长,最终可能 OOM。

生产建议:

java
new LinkedBlockingQueue<>(1000);

容量必须结合内存、任务大小、处理速度和可接受延迟来估算。

SynchronousQueue

SynchronousQueue 不存储元素,它是直接交接队列。

mermaid
flowchart TD
    A["生产者 put"] --> B{"是否有消费者正在 take"}
    B -->|有| C["直接交给消费者"]
    B -->|没有| D["生产者等待"]

特点:

特点说明
容量为 0不保存任务
直接交接生产者和消费者必须配对
常用于 cached 线程池没有空闲线程时倾向创建新线程

这解释了为什么 newCachedThreadPool 可能创建很多线程:它用 SynchronousQueue,不排队,任务交不出去就尝试创建新线程。

PriorityBlockingQueue

优先级阻塞队列,元素按优先级出队。

java
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 只有元素到期后才能被取出,常用于超时任务、延迟取消。

java
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 内存级别,服务重启会丢任务,分布式部署也难统一。

线程池和队列的关系

线程池提交任务的关键流程:

mermaid
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,还必须看队列类型和容量。

商业场景:异步通知发送

java
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);
    }
}

为什么 submitoffer 而不是 put

接口线程通常不能无限阻塞。队列满时返回 false,调用方可以降级、返回繁忙、记录失败或走 MQ。

商业场景:批量入库缓冲

java
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 或落库状态表。

线上排查:队列堆积

mermaid
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