Java 基础体系 · 第 59/100 篇。示例统一以 Java 25 LTS 为语言和 JVM 基线;框架示例使用与其兼容的现代稳定版本。

Java Queue 与 Deque:优先队列、双端队列、阻塞语义和选型

QueueDeque 都属于 Java 集合框架,但它们描述的不是某一种具体数据结构,而是对“元素从哪里进入、从哪里离开、在什么条件下操作失败”的抽象约束。

  • Queue<E> 通常表示单端进入、从队头取出的队列。
  • Deque<E>(double-ended queue)表示双端队列,队头和队尾都可以插入或删除。
  • PriorityQueue<E> 使用优先级决定队头,不是 FIFO 队列。
  • BlockingQueue<E>BlockingDeque<E> 在队列为空或已满时,可以等待而不是立即失败。
  • 并发队列还包括不阻塞的 ConcurrentLinkedQueueConcurrentLinkedDeque,以及用于任务传递的 SynchronousQueueLinkedTransferQueue

选择这些类型时,必须分别回答三个问题:

  1. 元素按什么规则离开:FIFO、优先级,还是两端都可操作?
  2. 容量不足或队列为空时:抛异常、返回特殊值,还是等待?
  3. 多线程访问时:是否需要线程安全、阻塞和可见性保证?

1. Queue 的抽象:队头、队尾与失败方式

Queue<E> 的核心概念是队头(head):

  • 删除操作从队头移除元素;
  • 查看操作查看队头但不移除;
  • 插入操作通常发生在队尾;
  • 队头由队列的排序规则决定。

对于普通 FIFO 队列,如果依次插入:

A -> B -> C

那么队头是 A,删除顺序是:

A -> B -> C

但是 Queue 并不强制所有实现都采用 FIFO。优先队列会根据优先级决定队头,延迟队列会根据延迟是否到期决定哪些元素可以被取出。

1.1 插入、删除和查看的两组方法

Queue 中有两组语义相近但失败行为不同的方法:

操作 失败时抛异常 失败时返回特殊值
插入 add(e) offer(e)
删除队头 remove() poll()
查看队头 element() peek()

例如:

Queue<String> queue = new ArrayDeque<>();

queue.add("A");
queue.offer("B");

System.out.println(queue.element()); // A
System.out.println(queue.peek());    // A

System.out.println(queue.remove());  // A
System.out.println(queue.poll());     // B

System.out.println(queue.poll());     // null

这里最后一次 poll() 返回 null,表示队列为空。

如果改用 remove()

Queue<String> empty = new ArrayDeque<>();
empty.remove(); // NoSuchElementException

如果改用 element()

Queue<String> empty = new ArrayDeque<>();
empty.element(); // NoSuchElementException

因此:

  • 能接受“当前没有元素”这一正常状态时,通常使用 poll()peek()
  • 逻辑上明确要求元素必须存在时,可以使用 remove()element()
  • 对有界队列,offer()add() 更适合把“容量已满”作为普通控制流处理。

1.2 null 不是可靠的队列空值协议

很多队列不允许插入 null,因为 poll()peek() 需要用 null 表示“没有元素”。

Queue<String> queue = new ArrayDeque<>();
queue.offer(null); // NullPointerException

Queue 接口本身允许某些实现接受 null,但官方文档明确指出,通常不推荐这样做,因为调用者无法区分:

queue.peek() == null

究竟表示:

  1. 队列为空;
  2. 队列中真实存在一个 null 元素。

LinkedList 作为 Queue 实现允许 null,但这种能力会让通用队列代码更容易产生歧义。并发队列和阻塞队列通常明确禁止 null


2. Deque:从两端操作的队列

Deque<E>Queue 的基础上增加了队头和队尾两个方向的操作:

方向 插入 删除 查看
队头 addFirst / offerFirst removeFirst / pollFirst getFirst / peekFirst
队尾 addLast / offerLast removeLast / pollLast getLast / peekLast

还有一组适合把 Deque 当作栈使用的别名:

栈语义 Deque 方法
入栈 push(e),等价于 addFirst(e)
出栈 pop(),等价于 removeFirst()
查看栈顶 peek(),通常查看队头

2.1 同一个 Deque 可以表达 FIFO 和 LIFO

作为 FIFO 队列:

Deque<String> fifo = new ArrayDeque<>();

fifo.addLast("A");
fifo.addLast("B");
fifo.addLast("C");

System.out.println(fifo.removeFirst()); // A
System.out.println(fifo.removeFirst()); // B

作为栈:

Deque<String> stack = new ArrayDeque<>();

stack.push("A");
stack.push("B");
stack.push("C");

System.out.println(stack.pop()); // C
System.out.println(stack.pop()); // B

两段代码使用的都是 ArrayDeque。区别不在数据结构本身,而在调用者选择的两端。

2.2 Deque 的双端操作不是“自动排序”

双端队列允许调用者分别操作两端,但它不会根据元素值自动排序:

Deque<Integer> deque = new ArrayDeque<>();
deque.addLast(10);
deque.addLast(2);
deque.addLast(7);

System.out.println(deque); // [10, 2, 7]

如果业务要求始终取出最小值或最高优先级元素,应使用 PriorityQueue,而不是试图通过 addFirstaddLast 模拟优先级。


3. PriorityQueue:队头是最小元素,不是整个集合有序

PriorityQueue<E> 是基于优先级的队列。默认情况下,队头是自然顺序最小的元素;如果提供了 Comparator,队头由比较器决定。

PriorityQueue<Integer> queue = new PriorityQueue<>();

queue.offer(7);
queue.offer(2);
queue.offer(5);

System.out.println(queue.peek()); // 2

while (!queue.isEmpty()) {
    System.out.println(queue.poll());
}

输出顺序为:

2
5
7

3.1 优先队列只保证队头,不保证迭代顺序

这是 PriorityQueue 最常见的误解:

PriorityQueue<Integer> queue = new PriorityQueue<>();
queue.add(7);
queue.add(2);
queue.add(5);
queue.add(1);

System.out.println(queue);

输出不应被当作排序结果。PriorityQueue 的迭代器和 toString() 不保证按优先级顺序遍历。

如果需要有序输出,必须反复调用 poll()

while (!queue.isEmpty()) {
    System.out.println(queue.poll());
}

原因在于,优先队列通常使用实现。最小堆只保证:

堆顶 <= 每个子节点

它不保证数组中任意两个位置都满足全局排序关系。

例如,以下数组可能表示一个最小堆:

[1, 2, 5, 7, 9, 6]

它满足:

  • 1 <= 21 <= 5
  • 2 <= 72 <= 9
  • 5 <= 6

但数组顺序并不是:

1, 2, 5, 6, 7, 9

3.2 堆操作的复杂度

对于通常的二叉堆实现:

  • peek():访问堆顶,时间复杂度通常为 O(1)
  • offer():插入后向上调整,通常为 O(log n)
  • poll():移除堆顶后向下调整,通常为 O(log n)
  • remove(Object):需要搜索目标,通常为 O(n)
  • contains(Object):通常为 O(n)
  • 遍历:不是有序遍历。

这里的 n 是当前元素数量。复杂度描述的是常见实现的算法性质;PriorityQueue 的接口契约主要规定行为,不要求调用者依赖某个具体内部实现。

3.3 比较器决定“谁优先”

默认优先级是自然顺序:

PriorityQueue<String> queue = new PriorityQueue<>();
queue.add("banana");
queue.add("apple");
queue.add("pear");

System.out.println(queue.poll()); // apple

如果希望数值越大优先级越高:

PriorityQueue<Integer> queue =
        new PriorityQueue<>(Comparator.reverseOrder());

queue.add(7);
queue.add(2);
queue.add(5);

System.out.println(queue.poll()); // 7

对象通常应显式定义比较器:

record Task(String name, int priority) {}

PriorityQueue<Task> tasks = new PriorityQueue<>(
        Comparator.comparingInt(Task::priority)
                  .reversed()
);

tasks.offer(new Task("低优先级任务", 1));
tasks.offer(new Task("高优先级任务", 10));
tasks.offer(new Task("中优先级任务", 5));

System.out.println(tasks.poll());
// Task[name=高优先级任务, priority=10]

比较器必须对所有可能入队的元素定义有效的排序关系。否则可能出现:

  • ClassCastException
  • 优先级顺序不符合业务预期;
  • 比较器与 equals 不一致导致去重逻辑产生误解;
  • 元素入队后用于排序的字段被修改,导致堆内部关系失效。

例如:

final class Job {
    int priority;
    Job(int priority) {
        this.priority = priority;
    }

    @Override
    public String toString() {
        return "Job(" + priority + ")";
    }
}

PriorityQueue<Job> queue =
        new PriorityQueue<>(Comparator.comparingInt(j -> j.priority));

Job job = new Job(1);
queue.offer(job);

job.priority = 100; // 不会自动重新调整堆

System.out.println(queue.peek()); // 仍可能是这个对象,但堆关系已不再按新值维护

如果优先级发生变化,应移除后重新插入,或者使用新的任务版本,而不是直接修改已入队对象的排序字段。

3.4 相同优先级不保证 FIFO

下面两个任务的优先级相同:

record Task(String name, int priority) {}

PriorityQueue<Task> queue = new PriorityQueue<>(
        Comparator.comparingInt(Task::priority)
);

queue.offer(new Task("first", 1));
queue.offer(new Task("second", 1));

poll() 先返回哪个,不应依赖插入顺序。要实现“优先级相同时按序号 FIFO”,应把序号加入比较键:

record OrderedTask(int priority, long sequence, String name) {}

AtomicLong sequence = new AtomicLong();

PriorityQueue<OrderedTask> queue = new PriorityQueue<>(
        Comparator.comparingInt(OrderedTask::priority)
                  .thenComparingLong(OrderedTask::sequence)
);

queue.offer(new OrderedTask(1, sequence.getAndIncrement(), "first"));
queue.offer(new OrderedTask(1, sequence.getAndIncrement(), "second"));

比较规则现在是:

  1. 优先级小的先出;
  2. 优先级相同时,序号小的先出。

这不是 PriorityQueue 自动提供的稳定排序,而是业务显式编码到比较器中的结果。


4. 普通队列与并发队列的边界

ArrayDequeLinkedListPriorityQueue 都不是线程安全容器。多个线程同时修改它们时,不能仅通过“每次操作看起来很简单”来获得安全性。

Queue<Integer> queue = new ArrayDeque<>();

如果一个线程执行 offer(),另一个线程同时执行 poll(),没有同步措施就可能产生数据竞争。问题不仅是元素丢失,还包括内部状态损坏、可见性不足和竞态条件。

4.1 ArrayDeque

ArrayDeque 是常用的非线程安全双端队列:

  • 支持两端插入和删除;
  • 通常比用 LinkedList 作为队列更适合栈和队列场景;
  • 不允许 null
  • 容量会动态扩展;
  • 不提供并发安全保证。

适合:

单线程 FIFO 队列
单线程栈
单线程滑动窗口
DFS/BFS 辅助容器

不适合:

多个线程直接共享并修改
需要阻塞等待
需要固定容量保护

4.2 LinkedList

LinkedList 同时实现 ListDequeQueue,因此可以作为队列使用:

Deque<String> deque = new LinkedList<>();

但它允许 null,并且节点对象和链表指针会带来额外内存开销。若只需要单线程的队列或双端队列,通常优先考虑 ArrayDeque;若需要按索引访问,LinkedList 也通常不是高效选择,因为索引访问需要遍历链表。

4.3 ConcurrentLinkedQueueConcurrentLinkedDeque

ConcurrentLinkedQueue 是线程安全的无界 FIFO 队列,ConcurrentLinkedDeque 是线程安全的无界双端队列。它们采用非阻塞并发算法,空队列或暂时无法完成操作时不会等待。

Queue<String> queue = new ConcurrentLinkedQueue<>();

queue.offer("A");
String value = queue.poll(); // 可能得到 null

“线程安全”不等于“所有复合操作自动安全”:

if (!queue.isEmpty()) {
    String value = queue.poll();
}

这两步之间,其他线程可能已经取走了元素。正确方式通常是直接使用 poll(),并处理它返回 null 的结果。

并发队列的迭代器一般是弱一致的:

  • 不会因为并发修改而抛出 ConcurrentModificationException
  • 可能看到迭代开始后发生的部分修改;
  • 不保证形成某个时间点的精确快照。

size() 在并发变化时也不适合作为严格同步依据。即使返回了某个数值,下一刻队列也可能已经变化。


5. 阻塞语义:不是“线程安全”的同义词

阻塞队列是指:当操作暂时无法完成时,调用线程可以等待条件成立。

BlockingQueue<E> 增加了三类重要操作:

操作 语义
put(e) 队列满时等待
take() 队列空时等待
offer(e, timeout, unit) 最多等待指定时间
poll(timeout, unit) 最多等待指定时间

普通 offer()poll() 仍然是立即返回的:

BlockingQueue<String> queue = new ArrayBlockingQueue<>(2);

boolean first = queue.offer("A"); // true
boolean second = queue.offer("B"); // true
boolean third = queue.offer("C"); // false,不等待

String value = queue.poll();       // A
String absent = queue.poll();      // 如果为空则 null

阻塞版本的行为不同:

queue.put("C"); // 队列满时阻塞,直到有空间
String value = queue.take(); // 队列空时阻塞,直到有元素

5.1 阻塞队列的状态转换

以容量为 2 的有界队列为例:

初始:[]
put(A) -> [A]
put(B) -> [A, B]
put(C) -> 阻塞
take() -> 返回 A,队列变为 [B]
解除 put(C) 的阻塞 -> [B, C]

put(C) 之所以能够继续,不是因为它超时或失败,而是因为另一个线程的 take() 改变了队列状态,使“剩余容量大于 0”这一条件成立。

消费者也有对应过程:

初始:[]
take() -> 阻塞
producer.put(A) -> [A]
解除 take() 的阻塞 -> 消费者得到 A,队列回到 []

5.2 中断是阻塞操作的重要退出路径

put()take()、带超时的 offer()poll() 都可能抛出 InterruptedException

try {
    String value = queue.take();
    process(value);
} catch (InterruptedException e) {
    Thread.currentThread().interrupt();
    return;
}

捕获中断后重新设置中断标记,是因为抛出 InterruptedException 通常会清除当前线程的中断状态。若直接吞掉异常,线程上层可能无法知道应当停止等待或退出。

生产者同样需要处理:

try {
    queue.put(task);
} catch (InterruptedException e) {
    Thread.currentThread().interrupt();
    cancel(task);
}

不能把中断简单写成:

try {
    queue.put(task);
} catch (InterruptedException ignored) {
}

这会丢失取消和关闭信号,线程可能继续运行,导致应用停止时无法及时结束。

5.3 阻塞不等于无限等待

offer(e, timeout, unit)poll(timeout, unit) 给等待设置了边界:

boolean accepted = queue.offer(
        task,
        500,
        TimeUnit.MILLISECONDS
);

if (!accepted) {
    rejectOrRetry(task);
}
Task task = queue.poll(1, TimeUnit.SECONDS);

if (task == null) {
    recordNoTask();
}

超时返回 falsenull 时,不能直接推断系统已经故障。它可能只是:

  • 当前确实没有消费者;
  • 生产速度暂时低;
  • 队列暂时满;
  • 线程调度延迟;
  • 业务处理时间超过预期。

是否重试、丢弃、降级或报警,取决于任务的幂等性和业务时限。

5.4 阻塞队列提供跨线程的内存可见性

BlockingQueue 的并发契约包含内存一致性保证:线程向队列中放入元素之前的操作,发生在另一个线程成功取出该元素之后的操作之前。

因此,下面的传递是有定义的:

record Message(String body) {}

class Producer implements Runnable {
    private final BlockingQueue<Message> queue;

    Producer(BlockingQueue<Message> queue) {
        this.queue = queue;
    }

    @Override
    public void run() {
        Message message = new Message("ready");
        try {
            queue.put(message);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
    }
}

class Consumer implements Runnable {
    private final BlockingQueue<Message> queue;

    Consumer(BlockingQueue<Message> queue) {
        this.queue = queue;
    }

    @Override
    public void run() {
        try {
            Message message = queue.take();
            System.out.println(message.body());
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
    }
}

消费者成功 take() 到消息后,可以依赖队列提供的发布关系读取消息对象在入队前已经完成的状态。但这不意味着消息对象之后可以被多个线程无同步地随意修改;发布安全和后续可变共享是两个问题。


6. BlockingQueue 的几种具体语义

6.1 ArrayBlockingQueue:固定容量与背压

BlockingQueue<Task> queue = new ArrayBlockingQueue<>(1000);

ArrayBlockingQueue 使用数组和固定容量:

  • 创建时必须指定容量;
  • 可以选择公平策略;
  • 队列满时,生产者可以阻塞或超时;
  • 队列空时,消费者可以阻塞或超时。

固定容量的直接效果是形成背压:当消费者处理不过来时,生产者最终会被限制。

生产速度 > 消费速度
        ↓
队列逐渐增长
        ↓
达到容量上限
        ↓
生产者阻塞、超时或拒绝

容量并不是越大越好。容量过大可能把处理延迟和内存压力隐藏起来,使任务长时间滞留;容量过小则可能让生产者频繁等待。容量应与单个任务大小、消费速率、可接受延迟和故障恢复策略一起确定。

公平策略会影响等待线程的调度倾向,但会带来额外同步成本;它不等于业务级别的严格全局公平,也不改变队列内部 FIFO 元素顺序。

6.2 LinkedBlockingQueue:可选容量的链式阻塞队列

BlockingQueue<Task> queue = new LinkedBlockingQueue<>(1000);

指定容量时,它可以作为有界队列使用。不指定容量时,容量上限接近 Integer.MAX_VALUE,这并不等于“无限容量”:

  • 内存仍然有限;
  • 每个链表节点有额外开销;
  • 生产速度持续超过消费速度时,最终可能造成内存压力。

如果系统需要明确的背压,通常应显式设置容量,而不是依赖默认上限。

6.3 PriorityBlockingQueue:可阻塞获取,但不是有界背压

PriorityBlockingQueue 是线程安全的优先队列:

BlockingQueue<Task> queue =
        new PriorityBlockingQueue<>(
                11,
                Comparator.comparingInt(Task::priority)
        );

它的关键语义是:

  • 队列为空时,take() 会等待;
  • 插入通常不会因为容量限制而等待,因为它是无界队列;
  • put() 不会提供“队列满时背压”的效果;
  • 同优先级元素不保证 FIFO;
  • 迭代器不保证优先级顺序。

因此,PriorityBlockingQueue 适合:

多个线程并发提交任务,消费者始终取当前最高优先级任务

但不适合单独承担:

限制生产者速度、保护内存、形成严格容量边界

如果优先任务队列也必须有容量上限,需要在应用层控制提交,或使用额外的信号量、拒绝策略和任务生命周期管理。

6.4 SynchronousQueue:没有内部存储的直接移交

SynchronousQueue 的容量不是 0 个“可暂存元素”,而是没有用于排队存储的缓冲区:

BlockingQueue<String> handoff = new SynchronousQueue<>();

Thread producer = Thread.startVirtualThread(() -> {
    try {
        handoff.put("task");
    } catch (InterruptedException e) {
        Thread.currentThread().interrupt();
    }
});

Thread consumer = Thread.startVirtualThread(() -> {
    try {
        System.out.println(handoff.take()); // task
    } catch (InterruptedException e) {
        Thread.currentThread().interrupt();
    }
});

put() 必须等待某个消费者接收,take() 也必须等待某个生产者交付。它适合直接移交任务,不适合吸收突发流量。

6.5 LinkedTransferQueue:区分“排队”和“直接移交”

LinkedTransferQueue 提供 transfer(e)

  • 如果已经有消费者等待,元素可以直接交给消费者;
  • 如果没有等待消费者,生产者会等待,直到元素被取走。

它还提供 tryTransfer(e),用于尝试立即移交而不等待。这个语义比普通 offer() 更接近“只有消费者已经准备好时才交付”。


7. BlockingDeque:双端操作加阻塞

BlockingDeque<E> 同时具备 Deque 的双端能力和 BlockingQueue 的等待能力。

它提供类似以下方法:

putFirst(e)
putLast(e)

takeFirst()
takeLast()

offerFirst(e, timeout, unit)
offerLast(e, timeout, unit)

pollFirst(timeout, unit)
pollLast(timeout, unit)

一个典型用途是工作窃取模型的简化版本:

  • 工作线程优先从自己的队头取任务;
  • 其他线程从队尾窃取任务;
  • 队列为空时,工作线程等待。
BlockingDeque<Task> deque = new LinkedBlockingDeque<>(100);

deque.putFirst(highPriorityTask);
deque.putLast(normalTask);

Task local = deque.takeFirst();
Task stolen = deque.takeLast();

需要注意,BlockingDeque 本身不会自动实现完整的工作窃取调度协议。它只提供双端并发容器和阻塞操作;任务归属、窃取条件、关闭通知和重复执行防护仍由应用或执行框架负责。


8. 端到端示例:有界生产者—消费者流水线

下面的程序展示一个完整的生产者—消费者流程:

  • 生产者生成任务;
  • ArrayBlockingQueue 限制缓冲容量;
  • 消费者使用 take() 等待任务;
  • 使用特殊的终止对象通知消费者退出;
  • 统一处理线程中断。
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.BlockingQueue;

public class BlockingPipelineDemo {
    private static final int CAPACITY = 2;

    private record Task(int id) {}

    // 终止标记使用独立对象,不使用 null
    private static final Task STOP = new Task(-1);

    public static void main(String[] args) throws InterruptedException {
        BlockingQueue<Task> queue = new ArrayBlockingQueue<>(CAPACITY);

        Thread consumer = Thread.startVirtualThread(() -> {
            try {
                while (true) {
                    Task task = queue.take();

                    if (task == STOP) {
                        break;
                    }

                    System.out.println(
                            "consume task " + task.id()
                    );
                }
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
                System.out.println("consumer interrupted");
            }
        });

        Thread producer = Thread.startVirtualThread(() -> {
            try {
                for (int i = 1; i <= 5; i++) {
                    Task task = new Task(i);
                    queue.put(task);
                    System.out.println("produce task " + task.id());
                }

                queue.put(STOP);
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
                System.out.println("producer interrupted");
            }
        });

        producer.join();
        consumer.join();
    }
}

8.1 前置条件

该代码需要 Java 25 环境,并使用了 Java 的虚拟线程 API。虚拟线程不改变 BlockingQueue 的队列语义:

  • take() 仍然在队列为空时等待;
  • put() 仍然在队列已满时等待;
  • 中断仍然是退出等待的重要机制。

虚拟线程只是改变线程资源模型,不会把阻塞队列变成非阻塞结构。

8.2 数据流和容量效果

容量为 2 时,最多同时缓存两个尚未消费的任务:

producer -> put(task) -> [task, task] -> consumer.take()

如果消费者处理较慢,第三个尚未被取走的任务会让生产者在 put() 处等待。等待会在消费者取走一个元素后解除。

程序输出中的生产和消费交错顺序不应写死。例如,可能看到:

produce task 1
produce task 2
consume task 1
produce task 3
consume task 2
...

也可能因为线程调度而出现其他合法交错,但消费者从 FIFO 队列获取普通任务的顺序仍是 1, 2, 3, 4, 5

8.3 终止标记的边界

STOP 是一个特殊任务对象。使用对象身份判断:

if (task == STOP)

比使用普通任务字段更不容易与业务数据冲突。如果有多个消费者,通常需要为每个消费者放入一个终止标记,或者由一个消费者取到终止标记后重新放回队列:

消费者数量 = N
终止标记数量通常也应按退出协议设计为 N

不能把 null 当作终止标记,因为阻塞队列禁止 null,而且 null 也无法区分空队列和控制消息。


9. 任务优先级示例:为什么需要显式定义稳定顺序

下面实现“数字越大优先级越高;优先级相同时按提交序号先进先出”:

import java.util.Comparator;
import java.util.concurrent.PriorityBlockingQueue;
import java.util.concurrent.atomic.AtomicLong;

public class PriorityQueueDemo {
    private record Task(
            int priority,
            long sequence,
            String name
    ) {}

    public static void main(String[] args) {
        AtomicLong sequence = new AtomicLong();

        PriorityBlockingQueue<Task> queue =
                new PriorityBlockingQueue<>(
                        11,
                        Comparator.comparingInt(Task::priority)
                                  .reversed()
                                  .thenComparingLong(Task::sequence)
                );

        queue.offer(new Task(
                1, sequence.getAndIncrement(), "low"
        ));
        queue.offer(new Task(
                10, sequence.getAndIncrement(), "high-1"
        ));
        queue.offer(new Task(
                10, sequence.getAndIncrement(), "high-2"
        ));
        queue.offer(new Task(
                5, sequence.getAndIncrement(), "medium"
        ));

        while (!queue.isEmpty()) {
            System.out.println(queue.poll());
        }
    }
}

预期顺序是:

Task[priority=10, sequence=1, name=high-1]
Task[priority=10, sequence=2, name=high-2]
Task[priority=5, sequence=3, name=medium]
Task[priority=1, sequence=0, name=low]

这里有两个容易忽略的点:

  1. PriorityBlockingQueue 的线程安全只保护并发容器操作,不会自动为任务优先级变化重新排序。
  2. 如果多个生产者并发生成序号,AtomicLong 保证每次递增操作具有原子性;但“序号先后”代表的是获取序号的顺序,不一定等同于业务事件发生的真实时间顺序。

如果业务要求更复杂的调度规则,例如租户配额、截止时间、重试次数和优先级组合,应该把规则集中放在比较器或调度层中,避免让调用者通过修改已入队对象来改变排序。


10. QueueDeque、优先队列和阻塞队列的关系

可以把接口关系理解为能力叠加,而不是实现分类:

Queue
├── 普通单端队列语义
├── PriorityQueue:改变队头排序规则
└── BlockingQueue:增加等待、超时和并发语义

Deque
├── Queue 的队列能力
├── 增加队头/队尾操作
└── BlockingDeque:增加双端阻塞语义

但接口继承关系和实际行为仍需区分:

  • Deque 继承 Queue,因此可以使用 offerpollpeek
  • PriorityQueue 实现 Queue,但不实现 Deque
  • ArrayDeque 实现 Deque,但不具备优先级排序;
  • PriorityBlockingQueue 实现 BlockingQueue,不是双端队列;
  • LinkedBlockingDeque 实现 BlockingDeque,但不会自动按优先级排序。

11. 选型:先确定排序,再确定等待和容量

11.1 只需要单线程 FIFO

选择:

Queue<E> queue = new ArrayDeque<>();

它适合任务暂存、BFS、事件顺序处理等场景。若需要可观察的固定容量或显式失败,可以考虑其他结构或在外层维护容量约束。

11.2 需要单线程栈或双端操作

选择:

Deque<E> deque = new ArrayDeque<>();
  • FIFO:addLast + removeFirst
  • LIFO:push + pop
  • 两端窗口:offerFirstpollLast

11.3 需要按优先级取出

选择:

Queue<E> queue = new PriorityQueue<>(comparator);

如果多个线程并发访问并且需要空队列等待:

BlockingQueue<E> queue =
        new PriorityBlockingQueue<>(11, comparator);

但要记住,PriorityBlockingQueue 默认不能提供有界背压。

11.4 需要多个线程共享 FIFO 队列,但不希望等待

选择:

Queue<E> queue = new ConcurrentLinkedQueue<>();

典型操作是:

E value = queue.poll();
if (value != null) {
    process(value);
}

这里的线程安全不会让 poll() 等待元素到达。

11.5 需要生产者—消费者等待

选择 BlockingQueue 的具体实现:

new ArrayBlockingQueue<>(capacity)

适合固定容量和背压。

new LinkedBlockingQueue<>(capacity)

适合链式结构和显式容量。

new LinkedBlockingQueue<>()

表示接近无界,但不应误解为没有内存风险。

new SynchronousQueue<>()

适合直接移交,不缓存。

new PriorityBlockingQueue<>(comparator)

适合优先级获取,但不提供容量上限。

11.6 需要双端并发和等待

选择:

BlockingDeque<E> deque =
        new LinkedBlockingDeque<>(capacity);

它适合双端任务调度、生产者根据任务类型选择插入端、消费者根据策略选择取出端等场景。


12. 常见误解与失败表现

12.1 把 PriorityQueue 当成排序容器

错误代码:

PriorityQueue<Integer> queue = new PriorityQueue<>();
queue.add(3);
queue.add(1);
queue.add(2);

for (Integer value : queue) {
    System.out.println(value);
}

这段代码没有排序输出保证。应该使用:

while (!queue.isEmpty()) {
    System.out.println(queue.poll());
}

如果需要保留原队列,先复制一份再消费副本。

12.2 用 isEmpty()poll() 组成并发判断

错误模式:

if (!queue.isEmpty()) {
    process(queue.poll());
}

在并发队列中,检查和取出不是一个原子操作。其他线程可能在两步之间取走元素,导致 poll() 返回 null

应改成:

E value = queue.poll();
if (value != null) {
    process(value);
}

12.3 认为 PriorityBlockingQueue 能限制内存

错误理解:

BlockingQueue<Task> queue =
        new PriorityBlockingQueue<>(100);

构造参数 100 通常是初始容量提示,不是容量上限。它不会让第 101 个元素自动阻塞。

如果生产者速率可以长期超过消费者速率,应使用真正有界的队列,或在优先级调度器外部增加明确的容量和拒绝策略。

12.4 以为 Deque 的两端操作天然线程安全

错误理解:

Deque<Task> deque = new ArrayDeque<>();

Deque 只是接口,不提供并发保证。即使所有操作都是 addFirstpollLast,多个线程同时访问普通 ArrayDeque 仍然是不安全的。

需要并发且不阻塞时使用:

Deque<Task> deque = new ConcurrentLinkedDeque<>();

需要并发且可以等待时使用:

BlockingDeque<Task> deque =
        new LinkedBlockingDeque<>(capacity);

12.5 在任务入队后修改排序字段

这种修改不会通知优先队列重新建立堆。表现可能是:

任务对象的 priority 字段已经改变,
但 poll() 返回顺序仍不符合新字段。

诊断时应检查:

  • 元素是否为可变对象;
  • 比较器读取的字段是否在入队后被修改;
  • 是否有同优先级顺序的额外规则;
  • 是否把迭代结果误当成优先级顺序。

13. 关闭、取消和故障路径

阻塞队列没有统一的 close() 方法。应用必须自行定义生产者和消费者的生命周期协议。

13.1 生产者失败

如果生产者在 put() 时被中断:

  1. put() 抛出 InterruptedException
  2. 当前线程的中断状态通常已被清除;
  3. 应恢复中断状态;
  4. 根据任务是否可重试决定取消、记录或转移任务。
catch (InterruptedException e) {
    Thread.currentThread().interrupt();
    cancelCurrentTask();
}

13.2 消费者失败

如果消费者处理任务时抛出运行时异常,任务可能已经从队列移除但尚未处理完成:

Task task = queue.take();
process(task); // 此处失败

队列不会自动把 task 放回去。是否重试必须由应用负责,例如:

Task task = queue.take();

try {
    process(task);
} catch (RuntimeException e) {
    retryOrDeadLetter(task, e);
}

把任务取出和业务处理视为一个不可失败的原子事务,是错误的假设。队列通常只负责传递,不负责业务处理结果的持久化和回滚。

13.3 停止顺序

一个有界生产者—消费者系统的正常停止通常需要明确顺序:

停止接收新任务
    ↓
生产者结束
    ↓
发送终止信号或触发取消
    ↓
消费者排空或放弃剩余任务
    ↓
等待线程退出

如果先停止消费者,生产者可能永久阻塞在满队列的 put() 上。若直接中断消费者而不处理正在执行的任务,还需要定义任务是否重试、是否进入死信队列以及是否允许重复执行。


14. 规范保证、实现行为与经验判断

使用队列时,应区分三类信息。

14.1 接口规范保证

这些是调用者可以依赖的语义:

  • Queue.poll() 在没有元素时返回 null
  • Queue.remove() 在没有元素时抛 NoSuchElementException
  • BlockingQueue.take() 在为空时等待;
  • 阻塞方法响应线程中断并抛出 InterruptedException
  • 优先队列的队头由自然顺序或比较器决定;
  • 优先队列迭代器不保证排序顺序;
  • BlockingQueue 不允许使用 null 作为元素。

14.2 常见实现特征

这些通常成立,但不应把具体内部实现当作接口契约:

  • PriorityQueue 使用堆;
  • ArrayDeque 使用可扩展数组;
  • LinkedBlockingQueue 使用链式节点;
  • 并发链队列使用非阻塞算法;
  • 某些公平策略会牺牲部分吞吐量。

如果代码依赖内部数组布局、节点结构或特定迭代顺序,就已经超出了集合接口的稳定抽象。

14.3 工程取舍

这些需要根据业务验证:

  • 队列容量应该多大;
  • 超时后是重试、丢弃还是降级;
  • 同优先级是否必须稳定排序;
  • 消费失败是否重试;
  • 关闭时是排空任务还是立即取消;
  • 是否需要持久化队列来应对进程崩溃。

集合队列通常只存在于当前 JVM 内存中。进程退出、机器故障或 JVM 崩溃后,未处理元素不会自动恢复。需要跨进程、持久化、确认和重放语义时,应使用具备这些能力的消息系统,而不能仅通过更换 Queue 实现解决。


15. 一个实用的决策表

需求 推荐类型 关键注意点
单线程 FIFO ArrayDeque 不允许 null,非线程安全
单线程栈 ArrayDeque 使用 push / pop
单线程双端操作 ArrayDeque 两端都可插入和删除
多线程、非阻塞 FIFO ConcurrentLinkedQueue poll() 不等待,迭代弱一致
多线程、非阻塞双端队列 ConcurrentLinkedDeque 无界,不提供等待
有界生产者—消费者 ArrayBlockingQueue 固定容量,能形成背压
可选容量的阻塞 FIFO LinkedBlockingQueue 默认上限很大,不等于无内存风险
阻塞优先级队列 PriorityBlockingQueue 无界,不提供容量背压
直接线程间移交 SynchronousQueue 不缓存元素
双端阻塞队列 LinkedBlockingDeque 支持两端等待、插入和删除
单线程优先级获取 PriorityQueue 迭代不排序,同优先级不保证 FIFO

最终选型不应从“哪个类最常见”开始,而应从状态和失败路径开始:

需要什么出队顺序?
    FIFO -> Queue
    优先级 -> PriorityQueue / PriorityBlockingQueue
    两端策略 -> Deque / BlockingDeque

需要等待吗?
    不等待 -> ArrayDeque / ConcurrentLinkedQueue 等
    空或满时等待 -> BlockingQueue / BlockingDeque

需要容量上限吗?
    需要 -> ArrayBlockingQueue、指定容量的 LinkedBlockingQueue/Deque
    不需要或直接移交 -> 相应无界队列或 SynchronousQueue

只要同时明确了排序规则、线程模型、空满状态和任务失败处理,QueueDeque 的选择就不再是类名记忆问题,而是对数据流和生命周期语义的直接建模。


系列导航与关联阅读

官方资料

本文依据 Java、Spring 与相关项目官方文档重新梳理;正文、示例与生产清单由 WR BLOG 编写。