Queue

本文源码以本机 JDK 21.0.4 为准

Queue 接口概述

前面说到,List 关注的是“按照位置组织元素”,Set 关注的是“保存唯一成员”。Queue 则把重点放在了另一个问题上,新元素怎样进入容器,下一个应该处理哪个元素。

普通排队场景中,先到的先处理,也就是 FIFO;优先队列则按照比较规则选择队首;双端队列还能从两端插入和取出,模拟 LIFO 栈。因此,Queue 不能简单理解为“所有实现都按照插入顺序遍历”。

Queue 继承 Collection,没有公开下标访问方法。它提供两组最常用的操作:

操作 失败时抛异常 正常失败时返回特殊值
添加元素 add(e) offer(e)
移除队首 remove() poll()
查看队首 element() peek()

容量受限时,add() 可能抛 IllegalStateException,offer() 返回 false;空队列上,remove()/element() 抛 NoSuchElementException,poll()/peek() 返回 null。

“返回特殊值”只是对这些正常失败状态的选择,不是所有错误都被吞掉。添加 null、传入不能比较的元素等,offer() 仍然可能抛出异常。

大多数这里介绍的 Queue 实现不接受 null,因为 null 已被用作空队列或无可用结果的标记。LinkedList 是允许 null 的例外,但作为 Queue 使用时,空队列与真实 null 元素会出现返回值歧义,应尽量避免。

BlockingQueue 与 TransferQueue

如果容器不仅要保存元素,还要协调生产者、消费者等待,就会进入 java.util.concurrent 中的 BlockingQueue。

操作 抛异常 立即返回特殊值 等待条件满足 等待至超时
入队 add() offer() put() offer(e, timeout, unit)
出队 remove() poll() take() poll(timeout, unit)
查看队首 element() peek() 没有对应方法 没有对应方法

这里的“立即返回”表示不等待容量或元素条件,不表示绝对不会等待内部互斥锁。put()/take() 等可中断方法则需要正确处理 InterruptedException;条件等待一般使用 while 循环,醒来之后重新检查状态。

BlockingQueue 提供安全发布关系:某线程将元素放入队列之前的操作,先行发生于另一线程取得该元素之后的操作。但它不会替业务协调元素交付之后仍在发生的字段修改,也没有统一的 close() 方法。

TransferQueue 在 BlockingQueue 基础上进一步区分“放进容器”和“交付给消费者”。transfer() 可以等待消费者取得元素,tryTransfer() 则尝试直接交付;注意,消费者取得元素,不等于消费者已经完成业务处理。

Queue 与 Deque 的操作及等待语义

下面逐个看这些集合怎样把接口语义落到真实结构上。

PriorityQueue

基本特性

前面概述 Queue 时提到,队首不一定是最早加入的元素。PriorityQueue 就是这个例子:它按照优先级决定下一个应该取出的元素。

它继承 AbstractQueue,实现 Serializable,底层是数组表示的二叉堆。没有比较器时使用自然排序,有比较器时使用 Comparator 的规则;按照这套规则,堆顶是最小元素。

1
2
3
4
5
6
PriorityQueue<Integer> queue = new PriorityQueue<>();
Collections.addAll(queue, 30, 10, 20, 10);
System.out.println(queue.peek()); // 10
while (!queue.isEmpty()) {
System.out.println(queue.poll()); // 10、10、20、30
}

如果希望较大的数先出队,使用 Comparator.reverseOrder() 即可。这里的“最小”是比较规则下最靠前的元素,不一定是数值最小。

PriorityQueue 允许重复,比较结果为 0 也不会去重。这一点与 TreeSet 很不一样:TreeSet 用比较结果决定唯一成员,而 PriorityQueue 用比较结果决定堆位置。

对于优先级相同的元素,它不承诺先入先出。需要稳定顺序时,可以在比较器中增加单调序号作为第二排序条件。

它不接受 null,即使比较器能够处理 null,offer() 也会先做非空检查。没有比较器时,元素要能够相互比较;混入不兼容类型可能抛出 ClassCastException。

它不是线程安全的。多线程共享时需要外部同步,或者使用后面的 PriorityBlockingQueue;不过后者同样不提供同优先级的稳定出队顺序。

PriorityQueue 没有用户设置的固定容量上限,初始容量只控制最初的数组分配。容量不足会扩容,内存仍然有限,因此不能把“逻辑无界”理解成永远可以成功分配空间。

复杂度也要按操作区分:peek()/size() 为 O(1),堆调整的 offer()/poll() 为 O(log n),contains()/remove(Object) 因为先线性搜索,通常为 O(n)。某次 offer() 还可能触发 O(n) 的数组复制,连续操作下的摊还成本为 O(log n)。

结构分析

PriorityQueue 的核心字段包括 Object[] queue、int size、Comparator<? super E> comparator 和 int modCount。queue 保存元素引用,size 表示有效成员数,modCount 用于迭代检测。

无参构造最终创建初始长度为 11 的数组。与 ArrayList 的懒分配不同,这里会直接分配数组;显式初始容量必须至少为 1。

二叉堆并没有在数组外再创建左右子节点。对于下标 i:

1
2
3
父节点:(i - 1) >>> 1,适用于 i > 0
左孩子:(i << 1) + 1
右孩子:(i << 1) + 2
PriorityQueue 的数组与二叉堆对应关系

堆要求父元素不大于其孩子,保证根是最小值;但是不同子树之间没有整体排序。例如数组 [1, 4, 2, 5, 6, 3] 是合法小顶堆,却不是升序数组。

所以,遍历数组、iterator()、toArray() 或直接打印 PriorityQueue,都不能当作已经排序的结果。想按优先级依次取出,要反复 poll();想保留原队列,则先复制到新的 PriorityQueue 再 poll()。

从普通 Collection 构造时,通常先复制元素数组,再 heapify。来源是 SortedSet 时会继承它的比较器,已排序的数组可以直接满足堆性质;来源是 PriorityQueue 时也会保留排序规则,并按照实际来源类型选择复制或重新建堆路径。

因此,“把一个集合转成 PriorityQueue”并不总是“依次执行 n 次 offer()”。通用数组建堆从最后一个非叶节点向根执行下沉,时间为 O(n)。底部节点多、但需要下沉的高度低,所有层的工作量相加并不是 O(n log n)。

基本操作

添加元素:追加到堆底,然后上浮
1
2
3
4
5
6
7
8
9
10
11
12
// 先保证数组空间,再沿父节点调整新元素的位置。
public boolean offer(E e) {
if (e == null)
throw new NullPointerException();
modCount++;
int i = size;
if (i >= queue.length)
grow(i + 1);
siftUp(i, e);
size = i + 1;
return true;
}

i 是原 size,也就是新元素最初可以占据的堆底位置。数组不足则 grow(),随后 siftUp(),最后更新 size。add() 直接调用 offer()。

有比较器时使用 siftUpUsingComparator,无比较器时使用 siftUpComparable。自然排序路径如下:

1
2
3
4
5
6
7
8
9
10
11
12
13
// 不必反复交换两个槽位,先向下搬父元素,最后填入新元素。
private static <T> void siftUpComparable(int k, T x, Object[] es) {
Comparable<? super T> key = (Comparable<? super T>) x;
while (k > 0) {
int parent = (k - 1) >>> 1;
Object e = es[parent];
if (key.compareTo((T) e) >= 0)
break;
es[k] = e;
k = parent;
}
es[k] = key;
}

如果新元素小于父节点,父元素被搬到当前空位,新元素继续向上找;直到父节点不大于它,或者已经到根,再填入。

这里保存的是一个“待填入元素”与不断向上的空位。相较每一层都交换两个元素,这样可以减少写数组的次数,但逻辑上仍是堆上浮。

对于自然排序,方法一开始会把元素转成 Comparable。即使第一次添加时数组里没有其他元素,不具备 Comparable 能力的对象也不能因为“没有比较对象”就安全存入。

比较器最好使用 Integer.compare、Long.compare 等方法,不要直接返回 a.priority - b.priority,因为相减可能溢出,破坏比较方向。

扩容:小数组与大数组采用不同增长幅度
1
2
3
4
5
6
7
8
9
10
11
// 首选增量随旧容量变化,最终长度还要满足最小增长和溢出检查。
private void grow(int minCapacity) {
int oldCapacity = queue.length;
// 小数组翻倍,否则增长 50%
int newCapacity = ArraysSupport.newLength(oldCapacity,
minCapacity - oldCapacity,
oldCapacity < 64 ? oldCapacity + 2 : oldCapacity >> 1
);
// 复制数组引用槽位,元素对象本身没有被深拷贝。
queue = Arrays.copyOf(queue, newCapacity);
}

这里的 oldCapacity < 64 ? oldCapacity + 2 : oldCapacity >> 1 是首选增加量,不是新容量本身。

所以在没有其他边界影响时,旧容量小于 64,新容量约为 2 * oldCapacity + 2;否则约为旧容量的 1.5 倍。比如默认长度 11 扩容时,通常变成 24,而不是变成 13。

ArraysSupport.newLength 还会处理最小增长要求与大数组边界。不能把它直接简化成“无条件乘 1.5”,也不能把某次观察到的容量增长当作接口的永久保证。

扩容复制的是元素引用,不会复制每个业务对象。数组只是容量变大,元素的优先级规则不因此改变。

查看与取出队首:用最后一个元素填补根

peek() 直接读 queue[0],空队列返回 null。poll() 则需要恢复堆性质:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
// 保存原堆顶,把最后一个成员取出作为待下沉元素,并清空旧尾槽位。
public E poll() {
final Object[] es;
final E result;

if ((result = (E) ((es = queue)[0])) != null) {
modCount++;
final int n;
final E x = (E) es[(n = --size)];
es[n] = null;
if (n > 0) {
final Comparator<? super E> cmp;
if ((cmp = comparator) == null)
siftDownComparable(0, x, es, n);
else
siftDownUsingComparator(0, x, es, n, cmp);
}
}
return result;
}

原堆顶被取出后,最后一个元素不能简单塞在根上就结束。它可能大于某个孩子,需要沿较小孩子的方向逐层下沉。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
// 每次选择较小的孩子向上补位,直到待填元素可以停下。
private static <T> void siftDownComparable(int k, T x, Object[] es, int n) {
Comparable<? super T> key = (Comparable<? super T>)x;
int half = n >>> 1;
while (k < half) {
int child = (k << 1) + 1;
Object c = es[child];
int right = child + 1;
if (right < n &&
((Comparable<? super T>) c).compareTo((T) es[right]) > 0)
c = es[child = right];
if (key.compareTo((T) c) <= 0)
break;
es[k] = c;
k = child;
}
es[k] = key;
}

n >>> 1 是第一个叶节点的下标边界,只有前面的节点需要检查孩子。先假设左孩子更小,再比较右孩子;选出较小者后,如果待填元素已经不大于它,就可以停止。

为什么要清空原尾槽位?因为 size 已经减一,那个槽位不再属于队列,继续保留引用会使数组无必要地持有对象。清空引用帮助对象具备回收条件,不意味着此处同步触发垃圾回收。

element()/remove() 采用 AbstractQueue 中的异常式入口:先 peek()/poll(),如果得到 null,就抛 NoSuchElementException。它们不会因为名字与 Collection.remove(Object) 接近就具有相同含义。

查找与指定删除:比较顺序和 equals 各管什么

contains(Object) 和 remove(Object) 使用 equals() 线性寻找目标,不是根据比较器一路定位。

例如比较器只比较任务的优先级,两项任务优先级相同但 equals() 不同,它们可以同时入队。remove(task) 找到的是 equals() 匹配的成员,不是“任意比较结果为 0 的成员”。

找到数组下标后,removeAt() 用最后一个成员填补空位,先尝试下沉;如果仍停在原位置,再尝试上浮:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
// 指定位置的删除既可能需要下沉,也可能需要上浮;返回值供迭代器修正遍历。
E removeAt(int i) {
// assert i >= 0 && i < size;
final Object[] es = queue;
modCount++;
int s = --size;
if (s == i) // 移除的是最后一个元素
es[i] = null;
else {
E moved = (E) es[s];
es[s] = null;
siftDown(i, moved);
if (es[i] == moved) {
siftUp(i, moved);
if (es[i] != moved)
return moved;
}
}
return null;
}

为什么先下沉,再上浮?被移动的尾元素与删除位置之间不一定只存在一个方向的违反:它可能比孩子大,也可能比父节点小。下沉没有移动时,再检查上浮,就能恢复堆性质。

初始搜索为 O(n),后续堆修复为 O(log n),完整 remove(Object) 仍是 O(n)。即使已经拥有目标对象引用,也没有公开的“直接定位堆下标”接口。

使用扩展:优先级稳定性与 Top K

如果业务要求优先级相同的任务按先后处理,可以用不可变数据加顺序号:

1
2
3
4
5
6
record Task(int priority, long sequence, String name) {}
PriorityQueue<Task> tasks = new PriorityQueue<>(
Comparator.comparingInt(Task::priority).thenComparingLong(Task::sequence));
tasks.offer(new Task(1, 2, "B"));
tasks.offer(new Task(1, 1, "A"));
System.out.println(tasks.poll().name()); // A

顺序号必须确实代表业务需要的先后关系。多线程提交时,分配序号的时刻和实际成功入队的时刻也可能不同,不要自动把两个顺序混为一谈。

Top K 则利用容量为 K 的小顶堆:保留最大的 K 个值,每来一个新值,与当前最小保留值比较,必要时替换。对于 1 ≤ K ≤ n,时间 O(n log K),额外空间 O(K);结果数组本身仍不一定整体有序。

更重要的是,修改已经在堆中的元素的比较字段,不会自动重新建堆。堆位置反映的是此前入队和调整时的比较结果。需要改变优先级时,可以先移除再重新入队,或者使用不可变任务与新的任务对象。

迭代器

PriorityQueue.iterator() 返回内部 Itr。普通推进沿 queue 数组下标访问,不做排序,也不复制完整快照;next()/remove() 对 expectedModCount 做 Fail-Fast 检查。

它支持 iterator.remove(),但是这一步比普通数组列表复杂,因为 removeAt() 可能把尾元素上浮到已经遍历过的区域。

源码因此保存一个 ArrayDeque<E> forgetMeNot。如果 removeAt() 返回了向前移动的元素,迭代器会把它暂存进去,普通数组遍历结束后再补发,避免遗漏。

比如图中的 [1, 4, 2, 5, 6, 3],迭代器已经返回下标 3 的 5,随后删除它,末尾的 3 会上浮到下标 1。这个下标之前已经访问过,不能只沿原 cursor 继续;forgetMeNot 就记录这个 3。

相反,如果空位被后面元素补上、没有向前上浮,cursor 需要回退一次,重新访问补到刚删除位置的元素。两种修正对应 removeAt() 返回值是否为 null。

补发元素的删除还会使用 removeEq(),以对象身份匹配具体成员,避免 equals() 相等但实际是另一份成员的对象被删除。

源码定位:PriorityQueue.java:532–552。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
// 迭代删除根据 removeAt 的结果回退 cursor 或登记向前移动的尾成员。
public void remove() {
if (expectedModCount != modCount)
throw new ConcurrentModificationException();
if (lastRet != -1) {
E moved = PriorityQueue.this.removeAt(lastRet);
lastRet = -1;
if (moved == null)
cursor--;
else {
if (forgetMeNot == null)
forgetMeNot = new ArrayDeque<>();
forgetMeNot.add(moved);
}
} else if (lastRetElt != null) {
PriorityQueue.this.removeEq(lastRetElt);
lastRetElt = null;
} else {
throw new IllegalStateException();
}
// 保存或同步当前迭代器预期的结构修改次数。
expectedModCount = modCount;
}

lastRet 记录最近一次从数组返回的下标,lastRetElt 则记录从 forgetMeNot 补发的对象。前一种删除可以直接使用下标,后一种需要按身份在当前堆中寻找;删除后同步 expectedModCount,允许迭代器继续合法推进。

这份代码还解释了 IllegalStateException 的来源:如果当前没有合法的最近返回成员,就不能 remove()。它与发现外部结构修改而抛出的 ConcurrentModificationException 是两种不同错误。

这里能看出,iterator.remove() 的正确性来自迭代器与堆调整方法之间的协作,不是简单执行 queue.remove() 然后让 cursor 自己继续。

外部 offer()、poll()、remove() 等结构修改可能导致 ConcurrentModificationException;Fail-Fast 只是尽力发现错误使用,不能代替同步。

PriorityQueue 的 Spliterator 报告 SIZED、SUBSIZED、NONNULL,没有报告 SORTED 或 ORDERED。堆顶优先级正确,与整个迭代序列有排序契约,是两个不同层面的保证。

操作 主要成本 是否按优先级处理
peek() O(1) 返回最优先的堆顶
offer() 摊还 O(log n),扩容时单次可为 O(n) 维护堆
poll() O(log n) 取出最优先的堆顶
contains()/remove(Object) O(n) 先按 equals() 搜索
Collection 建堆 通常 O(n) 建立堆性质
iterator() 完整遍历 通常 O(n),迭代删除另有修复成本 不保证按优先级排列

ArrayBlockingQueue

基本特性

Java 阻塞队列的历史可以追溯到 JDK 1.5,当时增加了 java.util.concurrent,也就是常说的 JUC 包。ArrayBlockingQueue 是其中非常典型的生产者、消费者协调容器。

它继承 AbstractQueue,实现 BlockingQueue 和 Serializable,底层是固定长度的数组。构造时必须指定大于 0 的容量,此后不会自动扩容。

一句话概括就是:以循环数组保存 FIFO 成员,并在满和空时提供等待机制的有界队列。

1
2
3
4
5
6
ArrayBlockingQueue<String> queue = new ArrayBlockingQueue<>(2);
System.out.println(queue.offer("A")); // true
System.out.println(queue.offer("B")); // true
System.out.println(queue.offer("C")); // false,已经满了
System.out.println(queue.poll()); // A
System.out.println(queue.offer("C")); // true,有空位了

它允许重复元素,不允许 null,成员保持队列的 FIFO 顺序。这里的 FIFO 描述成功入队成员的出队顺序,不等于多个线程最早调用 put() 的线程一定最早完成。

它是线程安全的,生产和消费共享一把 ReentrantLock。构造时还能选择公平或非公平锁,默认非公平。公平策略关注等待线程获取访问机会的顺序,通常减少长期饥饿风险,但可能增加调度开销;它不是另一个元素排序规则。

有界容量的实际意义,是限制队列里能积压的任务数量。生产速度超过消费速度时,put() 等待,offer() 失败,超时 offer() 在预算耗尽时返回 false,业务可以据此选择继续等待、降级或拒绝。

但是有界队列只限制已入队的成员数量,不会限制所有生产线程已经创建、正在本地等待提交的对象。完整的负载控制还要结合生产者并发度和任务生命周期。

结构分析

核心存储字段为 final Object[] items;takeIndex 指向下一次出队位置,putIndex 指向下一次写入位置,count 记录有效元素数量。

这里不是 ArrayList 那种“有效元素都从下标 0 连续排列”。两个下标在数组末端回绕,逻辑上的连续队列可以跨过数组边界。

ArrayBlockingQueue 的循环数组、计数与条件队列

对于固定长度 5 的数组,即使 takeIndex 和 putIndex 恰好相同,也可能是空,也可能是满,必须结合 count 判断。ArrayBlockingQueue 使用 count 区分满空,可以使用全部容量,不需要永远空出一个槽位。

这一点与后面的 ArrayDeque 不同。两个类都用循环数组,但满空判定机制和容量策略不能直接混用。

锁与条件对象如下:lock 保护数组、两个下标和 count;notEmpty 管理等待取得元素的线程;notFull 管理等待空位的线程。两种 Condition 来自同一把锁。

源码定位:ArrayBlockingQueue.java:270–277。

1
2
3
4
5
6
7
8
9
// 容量确定后分配数组,生产和消费的两个条件都绑定到同一把锁。
public ArrayBlockingQueue(int capacity, boolean fair) {
if (capacity <= 0)
throw new IllegalArgumentException();
this.items = new Object[capacity];
lock = new ReentrantLock(fair);
notEmpty = lock.newCondition();
notFull = lock.newCondition();
}

因此,数组首尾虽然是不同位置,但生产与消费的结构操作不能像 LinkedBlockingQueue 那样分别由两把端点锁并行执行。

为什么仍然可能有很好的吞吐?数组不必为每一个元素创建链表节点,首尾操作也只修改固定数量的槽位与计数;真正耗时还取决于锁竞争、任务处理和调度,不能只从“单锁”得出性能一定差的结论。

基本操作

入队:写槽位、移动下标、通知消费者

实际入队由私有 enqueue() 完成,调用它之前必须已经持有 lock,并确认存在空位。

源码定位:ArrayBlockingQueue.java:179–188。

1
2
3
4
5
6
7
8
9
10
11
// 把元素写入 putIndex,回绕下标后增加 count,并通知等待非空的线程。
private void enqueue(E e) {
// assert lock.isHeldByCurrentThread();
// assert lock.getHoldCount() == 1;
// assert items[putIndex] == null;
final Object[] items = this.items;
items[putIndex] = e;
if (++putIndex == items.length) putIndex = 0;
count++;
notEmpty.signal();
}

putIndex 到达数组长度时回到 0。它改变的是数组位置,不会搬动其他现有成员,因此不考虑等待的正常尾部入队结构成本为 O(1)。

enqueue() 末尾调用 notEmpty.signal()。signal() 并不会让消费者立刻在生产者仍持锁时读取数组,它只是使一个等待线程具备后续竞争锁的机会;消费者真正继续执行,要先重新获得锁。

普通 offer() 先检查 null,再 lock.lock(),在锁内检查 count 是否已经等于 items.length。满则返回 false,不满则 enqueue()。

注意,offer() 不等容量,不表示 lock.lock() 永不等待。在其他线程持有锁时,普通 offer() 仍可能短暂等待锁。这与 ConcurrentLinkedQueue 那种 CAS 为主的更新路径不同。

put:为什么条件等待要写 while

源码定位:ArrayBlockingQueue.java:364–375。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
// 获取锁和等待空位均可被中断;被唤醒后要再次确认队列没有满。
public void put(E e) throws InterruptedException {
Objects.requireNonNull(e);
final ReentrantLock lock = this.lock;
// 获取锁时也响应线程中断。
lock.lockInterruptibly();
try {
while (count == items.length)
// 队列满时释放锁等待,返回后重新检查条件。
notFull.await();
enqueue(e);
} finally {
lock.unlock();
}
}

这里使用 lockInterruptibly(),允许线程在等待获取锁时响应中断。满队列上的 notFull.await() 会释放 lock,让消费者有机会取得元素;被通知后,当前线程先重新获得锁,再检查条件。

为什么使用 while 而不是 if?一方面 Condition 允许虚假唤醒;另一方面,即使通知来自真实出队动作,当前线程重新获得锁之前,其他生产者也可能先占用了新空位。

因此,通知只意味着“可以重新检查”,不意味着“你已经预留了一个空位”。这个规律同样适用于后面其他阻塞队列的空、满判断。

出队:释放槽位,通知生产者

源码定位:ArrayBlockingQueue.java:194–208。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
// 保存队首元素并清空引用槽位,移动 takeIndex,释放一个容量单位。
private E dequeue() {
// assert lock.isHeldByCurrentThread();
// assert lock.getHoldCount() == 1;
// assert items[takeIndex] != null;
final Object[] items = this.items;
@SuppressWarnings("unchecked")
E e = (E) items[takeIndex];
items[takeIndex] = null;
if (++takeIndex == items.length) takeIndex = 0;
count--;
if (itrs != null)
itrs.elementDequeued();
notFull.signal();
return e;
}

dequeue() 清空旧槽位,使数组不再无必要地保留已取出元素。count 减一后通知 notFull,等待生产者可以重新竞争锁并检查空位。

take() 与 put() 对称:

源码定位:ArrayBlockingQueue.java:415–425。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
// 空队列上等待 notEmpty,获得元素后执行同一套 dequeue。
public E take() throws InterruptedException {
final ReentrantLock lock = this.lock;
// 获取锁时也响应线程中断。
lock.lockInterruptibly();
try {
while (count == 0)
// 暂时没有可取元素,释放锁等待生产者通知。
notEmpty.await();
return dequeue();
} finally {
lock.unlock();
}
}

普通 poll() 则在锁内检查 count,空队列直接返回 null;peek() 在同一把锁下读取 takeIndex 的成员,不删除。

所有合法线程访问都通过锁和相应方法进行,不能自己拿到内部数组并绕过它修改。线程安全的意义是共享结构被协调,不是数组的每一个槽位都声明 volatile。

超时与中断

超时 offer() 将 timeout 转成纳秒,在满队列循环中调用 notFull.awaitNanos(),将返回值作为剩余时间。每次被唤醒后,既检查容量条件,也检查剩余预算。

1
2
3
4
5
ArrayBlockingQueue<String> queue = new ArrayBlockingQueue<>(1);
queue.put("A");
System.out.println(queue.offer("B", 20, TimeUnit.MILLISECONDS)); // false
System.out.println(queue.take()); // A
System.out.println(queue.poll(20, TimeUnit.MILLISECONDS)); // null

超时参数不是硬实时调度保证。线程可能等待重新获取锁,实际返回时刻也受到调度影响;不能仅凭 20ms 就要求方法恰好在第 20ms 返回。

InterruptedException 则表示等待被取消。调用链能够继续传播时,就继续 throws;需要在上层结束任务时,可以保留中断状态并退出,不要捕获后无条件继续原来的无限循环。

BlockingQueue 没有通用 close()。消费者什么时候停止,需要业务自行制定协议,例如独立的结束标记或协调中断。null 不能作为毒丸,因为队列拒绝 null;多个消费者也不能假定一个结束标记就能让所有线程退出。

按值删除、批量移出与容量观察

remove(Object) 需要在锁内线性扫描。删除队首可以直接推进 takeIndex;删除中间成员则要移动后续数组引用,把队列逻辑顺序接起来,同时更新 putIndex 与 count。

所以首尾入队出队 O(1),不意味着任意位置删除也 O(1)。contains()、指定成员删除等通常为 O(n),并在执行期间占用同一把锁。

drainTo() 可以批量把现有成员转移到另一个 Collection,最多移出指定数量。它在队列锁保护下遍历一段现有成员,减少多次加锁的成本,最后统一更新 count 和 takeIndex,并通知等待生产者。

但它不是涉及两个任意容器的通用事务。如果目标 Collection.add() 中途抛出异常,已经成功转移的部分不会自动全部回滚;目标不能是当前队列自身。

1
2
3
4
5
6
ArrayBlockingQueue<Integer> queue = new ArrayBlockingQueue<>(4);
queue.addAll(List.of(1, 2, 3));
List<Integer> batch = new ArrayList<>();
System.out.println(queue.drainTo(batch, 2)); // 2
System.out.println(batch); // [1, 2]
System.out.println(queue.peek()); // 3

size() 与 remainingCapacity() 在锁内读取 count,主要结构成本为 O(1)。但是释放锁之后状态就可能变化,if (remainingCapacity() > 0) add(e) 不能保证随后添加成功。真正决定能否接收成员的,应是实际 offer()/put() 的结果或行为。

addAll() 也不等于一次为整批预留容量。它沿通用添加路径逐个操作,队列容量不足可能已经接受部分元素再抛异常;业务需要整批原子提交时,应该把一批任务作为一个成员或另行设计协调。

迭代器

ArrayBlockingQueue 的迭代器是弱一致的,没有对整个遍历持有同一把锁,也没有在创建时复制完整数组。

但是“弱一致”并不是“读取时完全不锁”。Itr 会在创建、推进和删除等步骤中使用队列锁,协调数组下标与当前状态,并提前缓存 nextItem。

外部出队之后,先前缓存的成员仍可能被 next() 返回,这是为了保持迭代器自己的推进语义。迭代也不会把新加入的每一个成员保证都观察到。

循环数组还带来一个特殊问题:数组下标会重复使用,原本的下标 2 可能后来属于另一轮绕回的成员。只保存一个 cursor,不足以正确处理外部删除和回绕。

源码为此使用 Itrs 管理活跃迭代器,通过 WeakReference 注册它们,记录 takeIndex 的回绕周期,并在 dequeue()、队列变空、中间删除等时机通知调整。弱引用避免队列为了维护迭代记录长期强持有已经不用的迭代器。

Itr 也有 detached 状态,已经不再需要追踪数组变动时可以脱离注册结构;nextItem、nextIndex、lastRet 等字段分别负责缓存元素、定位和删除状态。

下面直接看 next() 如何从缓存返回成员,并在锁下推进:

源码定位:ArrayBlockingQueue.java:1212–1238。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
// 先兑现已缓存成员,再结合回绕和删除通知调整索引并寻找下一项。
public E next() {
final E e = nextItem;
if (e == null)
throw new NoSuchElementException();
final ReentrantLock lock = ArrayBlockingQueue.this.lock;
// 先获取保护队列结构的锁。
lock.lock();
try {
if (!isDetached())
incorporateDequeues();
// assert nextIndex != NONE;
// assert lastItem == null;
lastRet = nextIndex;
final int cursor = this.cursor;
if (cursor >= 0) {
nextItem = itemAt(nextIndex = cursor);
// assert nextItem != null;
this.cursor = incCursor(cursor);
} else {
nextIndex = NONE;
nextItem = null;
if (lastRet == REMOVED) detach();
}
} finally {
lock.unlock();
}
return e;
}

incorporateDequeues() 会根据保存的回绕周期和 takeIndex,判断某些原下标是否已经被消费。原槽位即使现在又非 null,也可能属于后来加入的另一轮成员,所以不能只凭数组当前位置重新解释旧的迭代记录。

remaining 的概念也不能简单套在这里:队列在遍历期间持续变化,Itr 维护的是自己尚可追踪的路径和缓存,而不是一个创建时固定数组副本。到达需要脱离状态时,它保留必要的 lastItem 信息,供后续删除核对身份。

remove() 的内部检查与迭代器注册关系,正是为了避免“旧下标指向新任务,误把新任务删掉”。迭代器的支持需要队列在出队、回绕、清空和中间移位时共同通知,而不仅是 next() 中加一次锁。

它支持 iterator.remove()。删除时需要核对当前下标和保存的成员是否仍然对应同一次观察,不能把一个已经复用的数组槽位直接当成原成员删除。

1
2
3
4
5
6
7
ArrayBlockingQueue<Integer> queue = new ArrayBlockingQueue<>(4);
queue.addAll(List.of(1, 2, 3));
Iterator<Integer> it = queue.iterator();
while (it.hasNext()) {
if (it.next() == 2) it.remove();
}
System.out.println(new ArrayList<>(queue)); // 无其他线程修改时:[1, 3]

迭代器不依赖外部修改就抛 ConcurrentModificationException 来工作,但同一个迭代器自身的推进状态仍然通常应由一个线程使用。集合线程安全,不自动把任意 Iterator 对象变成可由多个线程同时推进的协议。

Spliterator 使用 CONCURRENT、ORDERED、NONNULL 等特征,表示遍历支持并发修改且遵循队列方向。它不提供创建时成员快照,也不能作为精确统计某一统一时刻全部成员的工具。

LinkedBlockingQueue

基本特性

ArrayBlockingQueue 用固定数组保存待处理成员,那么如果希望按实际入队数量分配节点,同时让生产和消费在不同端点上尽量独立进行,就可以看看 LinkedBlockingQueue。

它继承 AbstractQueue,实现 BlockingQueue 和 Serializable,底层是单向链表,维护 FIFO 顺序,允许重复,不接受 null,而且线程安全。

它可以指定容量,也可以不指定。无参构造使用 Integer.MAX_VALUE 作为逻辑上限,通常被称为近似无界队列;这不是“不需要内存”,也不是永远不会满。

1
2
3
4
5
6
LinkedBlockingQueue<String> queue = new LinkedBlockingQueue<>(2);
queue.put("A");
queue.put("B");
System.out.println(queue.offer("C")); // false
System.out.println(queue.take()); // A
System.out.println(queue.offer("C")); // true

它的一个重要特点是生产与消费使用两把锁:putLock 保护尾部添加,takeLock 保护头部移除。同一端上的竞争仍然需要排队,但一个生产者与一个消费者在合适状态下可以分别执行端点操作。

不要把“锁分离”理解成所有方法都只用一把端点锁。有些遍历、按值删除、复制和清空操作,需要同时取得两把锁。

LinkedBlockingQueue 没有公开的公平锁构造参数,内部创建的是默认非公平 ReentrantLock。它也不提供从队尾取元素的双端接口,虽然名字含 Linked,但这里不是 LinkedList 那种双向链表。

有界与无参构造的选择,直接影响积压行为。显式容量可以让过快的生产者等待或感知失败;无参构造会容纳更多待处理任务,但任务长期积压仍会增加内存占用和排队时延。

结构分析

核心字段包括 final capacity、AtomicInteger count、head、last,以及两把锁和对应的 Condition。

head 指向一个 item 为 null 的哨兵节点,真实的第一个成员位于 head.next;last 指向最后一个节点。空队列中,head 和 last 指向同一个哨兵。

LinkedBlockingQueue 的哨兵与两把锁

节点只保存 item 与 next。Node.next 还可能指向节点自己,表示这个旧节点已经离开当前链表,后续遍历要转回新的 head.next,而不是沿它原来的历史链路一直向前。

这里的“哨兵 item 为 null”不表示用户可以保存 null。它是内部结构标记,用来简化首尾边界操作。

无参构造没有提前创建 Integer.MAX_VALUE 个槽位,只创建当前需要的哨兵和锁等对象,真实节点随入队分配。数组队列的固定槽位成本与链表队列的逐节点分配成本,需要分开比较。

count 为什么要用 AtomicInteger?因为 putLock 与 takeLock 是不同的锁,生产与消费都会修改成员数。普通 int 不能仅靠两把互不相同的锁就自然成为安全共享计数。

这里需要一个能被两端共同观察并原子更新的计数器,而且该计数器还参与空、满判断和跨端通知。LongAdder 更适合分散统计累加,不能直接替换需要精确原子递增、递减与旧值判断的 count。

基本操作

尾部添加与首部摘取

入队的结构方法非常短:

源码定位:LinkedBlockingQueue.java:201–205。

1
2
3
4
5
6
// 尾部添加由 putLock 保护,连接新节点并移动 last。
private void enqueue(Node<E> node) {
// assert putLock.isHeldByCurrentThread();
// assert last.next == null;
last = last.next = node;
}

实际尾部指针调整为 O(1),不需要从 head 一路寻找末尾。添加前会创建一个 Node,保存业务元素引用;节点分配与对象处理成本需要另外考虑。

出队也使用固定数量的指针调整:

源码定位:LinkedBlockingQueue.java:212–222。

1
2
3
4
5
6
7
8
9
10
11
12
// 原首个数据节点转成新的哨兵,旧哨兵自连接以切断历史链路。
private E dequeue() {
// assert takeLock.isHeldByCurrentThread();
// assert head.item == null;
Node<E> h = head;
Node<E> first = h.next;
h.next = h; // 帮助 GC
head = first;
E x = first.item;
first.item = null;
return x;
}

first 原来是第一个数据节点。取出它的 item 后,这个节点被设置为新的 head,并清空 item;真正剩余的首成员成为它的 next。

旧哨兵 h 的 next 被设为 h 自己,而不是简单保留对新 head 的引用。迭代器可以识别这种自连接并回到当前链表,旧节点也不会通过长链继续无必要地持有后续节点。

put:生产端锁与计数发布

源码定位:LinkedBlockingQueue.java:326–354。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
// 生产端等待容量,链接节点后原子增加 count,再按状态变化通知相应等待者。
public void put(E e) throws InterruptedException {
if (e == null) throw new NullPointerException();
final int c;
final Node<E> node = new Node<E>(e);
final ReentrantLock putLock = this.putLock;
final AtomicInteger count = this.count;
putLock.lockInterruptibly();
try {
/* 注意:count 虽然不受 lock 保护,却仍用作等待条件。
* 这是可行的,因为此处 count 只可能减少(其他 put 都被 lock 挡在外面),
* 而一旦它从满容量发生变化,我们(或其他等待中的 put)会收到通知。
* 其他等待条件中使用 count 的方式同理。
*/
while (count.get() == capacity) {
// 队列满时释放锁等待,返回后重新检查条件。
notFull.await();
}
enqueue(node);
// 原子增加计数,返回增加之前的成员数量。
c = count.getAndIncrement();
if (c + 1 < capacity)
notFull.signal();
} finally {
putLock.unlock();
}
if (c == 0)
signalNotEmpty();
}

先持有 putLock,满时等待 notFull。由于其他生产者也受这把锁限制,在当前等待条件检查期间,计数可能由消费者减少;醒来之后仍必须重新检查。

enqueue() 完成后,count.getAndIncrement() 返回增加前的数量 c。这个旧值有两个用途。

如果 c + 1 < capacity,队列还有空位,可以继续 notFull.signal(),让另一个等待生产者获得机会;如果 c 为 0,这次添加使队列从空变为非空,需要通知消费者。

注意,源码不是在持有 putLock 时直接获取 takeLock,而是先释放生产端锁,再调用 signalNotEmpty()。后者短暂持有 takeLock,通知 notEmpty。

这样的顺序避免简单地在两端操作中形成相反的双锁嵌套。双锁协调仍然有严格次序,不是想到哪一把就取哪一把。

count 的原子更新和可见性,还与另一端观察已经链接的节点相配合。不能把节点的写入、计数增加、通知拆成随意重排的独立步骤。

take:消费端对称协调

源码定位:LinkedBlockingQueue.java:427–447。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
// 消费端等待成员,摘下队首后减少 count;原来满时再通知生产端。
public E take() throws InterruptedException {
final E x;
final int c;
final AtomicInteger count = this.count;
final ReentrantLock takeLock = this.takeLock;
takeLock.lockInterruptibly();
try {
while (count.get() == 0) {
// 暂时没有可取元素,释放锁等待生产者通知。
notEmpty.await();
}
x = dequeue();
// 原子减少计数,返回减少之前的成员数量。
c = count.getAndDecrement();
if (c > 1)
notEmpty.signal();
} finally {
takeLock.unlock();
}
if (c == capacity)
signalNotFull();
return x;
}

count 为 0 时等待 notEmpty;取出后,getAndDecrement() 返回减少前的数量 c。如果 c 大于 1,说明还有数据,可以通知另一个消费者;如果 c 等于 capacity,说明从满变为非满,需要在释放 takeLock 后通知生产者。

这里不要求每次入队都同时唤醒所有消费者,也不要求每次出队都同时唤醒所有生产者。利用状态边界与同端传递通知,可以减少不必要的唤醒竞争。

如果 count 更新了,但另一端线程还没有被调度,元素状态并不是“尚未生效”。通知负责使等待者继续工作,线程什么时候运行还受调度与锁竞争影响。

offer、poll 与超时重检

普通 offer() 在锁外先读取 count,发现已满可以提前返回。发现有空间后,还需要获取 putLock,再次检查容量,防止其他生产者已经先占用空位。

因此,外面的检查只是减少明显失败情况下的锁开销,不是最终的容量预留。不能把源码简化成“先读 count,小于容量就直接接节点”。

普通 poll() 也有类似优化:锁外发现 count 为 0 可以直接返回 null,否则获取 takeLock 后重新确认。peek() 需要在消费端协调下访问第一个有效成员。

超时 offer()/poll() 将等待预算转换成纳秒,循环等待对应 Condition。等待可能被中断,返回值表示这次操作是否在其观察到的条件下完成,不代表状态之后永远不变。

1
2
3
4
5
LinkedBlockingQueue<Integer> queue = new LinkedBlockingQueue<>(1);
System.out.println(queue.poll(10, TimeUnit.MILLISECONDS)); // null
System.out.println(queue.offer(1, 10, TimeUnit.MILLISECONDS)); // true
System.out.println(queue.remainingCapacity()); // 此刻为 0
System.out.println(queue.take()); // 1

size() 直接读取 AtomicInteger count,remainingCapacity() 为 capacity 减去该值,结构读取为 O(1)。多次读取并不形成统一状态快照,也不能用“先检查剩余容量”替代实际 offer() 的最终检查。

按值删除与完整锁定

remove(Object) 和 contains(Object) 需要遍历链表。为了不让遍历和结构删除同时受到两端改变影响,这些方法会通过 fullyLock() 同时取得 putLock 和 takeLock。

源码定位:LinkedBlockingQueue.java:227–230。

1
2
3
4
5
// 需要稳定协调两端时,统一先取得 putLock,再取得 takeLock。
void fullyLock() {
putLock.lock();
takeLock.lock();
}

释放则按相反顺序进行。这个锁顺序与跨端 signal() 先释放本端锁的设计一起构成协调约束。

所以“生产者与消费者可以并行”主要适用于普通端点路径。执行按值删除、完整数组复制或清空时,仍然可能阻塞两端,不能将所有方法都描述成互不影响。

remove() 按 equals() 搜索,整体通常为 O(n);已经定位节点后,跳过节点和修复 last 为 O(1),但不消除此前的搜索成本。

drainTo:只固定消费端,允许生产端继续

drainTo() 在 takeLock 下取一批现有节点,生产端仍可能接入新节点。方法根据开始时选择的移出数量,逐个加入目标 Collection,再更新链表头与计数。

这个设计与 ArrayBlockingQueue 的单锁批量移出不同。它减少消费端多次加锁,同时不必把整个操作都放在双锁内,生产者可以继续在尾部工作。

1
2
3
4
5
LinkedBlockingQueue<Integer> queue = new LinkedBlockingQueue<>(List.of(1, 2, 3));
List<Integer> batch = new ArrayList<>();
queue.drainTo(batch, 2);
System.out.println(batch); // [1, 2]
System.out.println(queue.poll()); // 3

目标 add() 中途抛异常时,已经转移部分需要按实际进度更新结构,不提供全部回滚。drainTo() 也不能传当前队列作为目标,更不能当作两个独立队列之间自动无死锁、可回滚的事务机制。

clear() 则需要完整锁定两端,断开或自连接旧节点,清空成员引用,并重置哨兵和 last;原来满时会通知生产端。

限时加入与双重容量检查

源码定位:LinkedBlockingQueue.java:365–390。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
// 带超时的生产者在 putLock 下重新检查容量,并以剩余预算等待 notFull。
public boolean offer(E e, long timeout, TimeUnit unit)
throws InterruptedException {

if (e == null) throw new NullPointerException();
long nanos = unit.toNanos(timeout);
final int c;
final ReentrantLock putLock = this.putLock;
final AtomicInteger count = this.count;
putLock.lockInterruptibly();
try {
while (count.get() == capacity) {
if (nanos <= 0L)
return false;
nanos = notFull.awaitNanos(nanos);
}
enqueue(new Node<E>(e));
// 原子增加计数,返回增加之前的成员数量。
c = count.getAndIncrement();
if (c + 1 < capacity)
notFull.signal();
} finally {
putLock.unlock();
}
if (c == 0)
signalNotEmpty();
return true;
}

外部或锁前看到 count 未满,只能作为快速观察。真正执行链接之前仍然必须在 putLock 下检查,防止多个生产者依据同一个剩余空位同时入队。

节点在进入等待前就已经创建,因此有界队列只限制链上有效成员,并不自动限制所有正在等待 put() 的生产线程已经分配的对象。控制内存时,生产者并发数量和对象大小仍需一起考虑。

等待完成后的 c 表示修改前数量:原来为 0 就需要跨到消费端通知 notEmpty;修改后仍有空位,则可以继续通知另一个生产者。通知顺序和条件都与容量边界有关,而不是每次无条件 signalAll()。

序列化同样需要协调稳定内容,写入容量与成员,再按读取结果重新建立节点。保存的是逻辑 FIFO 内容,不是原样保存每个 next 指针的地址。

迭代器

LinkedBlockingQueue 的 Itr 是弱一致迭代器,它没有在创建时复制完整链表,也不采用 modCount 的 Fail-Fast 策略。

Itr 保存 next、nextItem、lastRet,以及协助删除定位的 ancestor。nextItem 先缓存下一个待返回的元素,即使外部线程在随后出队,缓存仍可能被返回。

创建时短暂 fullyLock(),确定第一个成员;next() 推进时再次短暂 fullyLock(),寻找下一个 item 非 null 的节点。锁在方法结束时释放,不会从第一次 next() 到最后一次 next() 一直持有。

下面是源码中推进阶段的节选:

源码定位:LinkedBlockingQueue.java:789–806。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
// 推进时协调两端,跳过已经清空 item 的节点;这里不是整个遍历的快照复制。
public E next() {
Node<E> p;
if ((p = next) == null)
throw new NoSuchElementException();
lastRet = p;
E x = nextItem;
fullyLock();
try {
E e = null;
for (p = p.next; p != null && (e = p.item) == null; )
p = succ(p);
next = p;
nextItem = e;
} finally {
fullyUnlock();
}
return x;
}

遇到自连接旧节点时,succ() 会把它映射到当前 head.next,避免在旧哨兵上无限循环。已删除节点、尾部新节点与缓存成员如何被观察,都取决于当前推进位置。

它支持 iterator.remove(),按保存的物理节点定位。如果该节点已经被外部删除,源码会检查 item 状态,避免再次减少 count;否则在完整锁下执行 unlink() 并维护计数和必要通知。

1
2
3
4
5
6
LinkedBlockingQueue<Integer> queue = new LinkedBlockingQueue<>(List.of(1, 2, 3));
Iterator<Integer> it = queue.iterator();
while (it.hasNext()) {
if (it.next() == 2) it.remove();
}
System.out.println(new ArrayList<>(queue)); // 无其他修改时:[1, 3]

forEachRemaining() 还有批量读取路径:在锁内取一批成员引用,再在锁外调用用户 action,减少反复锁定及持锁执行回调的影响。这个分批机制仍然不是创建时的完整成员快照。

Spliterator 报告 ORDERED、NONNULL、CONCURRENT。CONCURRENT 描述并发修改下的遍历契约,不意味着 iterator.next() 没有锁,也不意味着整个遍历对应一个精确时刻。

把它与 ArrayBlockingQueue 比较时,可以这样看:数组队列容量固定、每个成员没有单独链表节点、普通生产消费共享一把锁;链表队列可指定上限、按成员分配节点、普通端点使用两把锁。选择还要结合容量控制、分配开销、遍历与删除频率,而不只记一句“双锁一定更快”。

PriorityBlockingQueue

基本特性

PriorityQueue 解决了优先处理问题,但是没有并发保护。PriorityBlockingQueue 则把优先级堆和阻塞式消费者等待结合起来。

它继承 AbstractQueue,实现 BlockingQueue 和 Serializable,允许重复元素,禁止 null,使用自然排序或 Comparator 规则决定堆顶。

它是逻辑无界的并发优先队列。构造时给出的 initialCapacity 是数组初始分配规模,不是容量上限;remainingCapacity() 返回 Integer.MAX_VALUE,也不是精确的“还可以分配多少元素”。

1
2
3
4
5
6
7
PriorityBlockingQueue<Integer> queue = new PriorityBlockingQueue<>(2);
queue.put(30);
queue.put(10);
queue.put(20); // 初始容量 2,不会因为已经有两个元素就等容量
System.out.println(queue.take()); // 10
System.out.println(queue.take()); // 20
System.out.println(queue.take()); // 30

由于没有普通容量满的等待条件,put() 和超时 offer() 不会为了等待空位而阻塞。但是它们仍然会竞争内部锁、可能分配数组,所以不能把这个行为表述成“方法绝对不等待任何资源”。

真正的阻塞主要发生在消费端:take() 在空队列上等待,超时 poll() 在没有可用成员时等待至条件满足或预算耗尽。

它不保证同优先级 FIFO。低优先级元素也可能因为持续到来的高优先级元素长期得不到处理,内部锁是否公平与业务优先级饥饿是两种不同的问题。

复杂度与 PriorityQueue 的堆操作接近:peek()/size() 的结构读取 O(1),offer()/take() 的堆调整 O(log n),指定值查询和删除通常 O(n)。实际还要加上锁竞争、数组分配及消费等待,不把等待时间写成算法上的固定常数。

结构分析

这个类没有简单地把普通 PriorityQueue 放在一把锁里当作所有操作的存储对象。它自己维护 Object[] queue、size、comparator,并实现上浮和下沉方法。

源码中的 PriorityQueue 字段 q 主要用于序列化协助,不应看到这个字段就推断所有读写都委派给 q。

主要协调对象是一把 ReentrantLock 和一个 notEmpty Condition。因为逻辑上不会因固定容量变满,正常入队不需要 notFull 条件队列。

PriorityBlockingQueue 的堆、单锁与扩容协调

还有一个 volatile allocationSpinLock,通过 VarHandle 协调新数组的分配。它与保护堆内容的主锁不是同一个职责:主锁保护成员与堆调整,分配标记尽量避免多个线程同时做冗余扩容分配。

这份源码默认初始容量也是 11。指定容量必须为正;来源集合构造则根据来源类型判断比较器、是否已经是兼容堆或排序成员,再选择复制和建堆路径。

堆在数组中的父子下标关系与 PriorityQueue 相同。比较器只保证每条父子关系符合优先级,不能让整个数组变成一张排好序的表。

基本操作

offer:在主锁中维护堆与成员数量

源码定位:PriorityBlockingQueue.java:463–484。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
// 拒绝 null,保证数组空间,在锁内上浮新成员,随后通知等待消费者。
public boolean offer(E e) {
if (e == null)
throw new NullPointerException();
final ReentrantLock lock = this.lock;
// 先获取保护队列结构的锁。
lock.lock();
int n, cap;
Object[] es;
while ((n = size) >= (cap = (es = queue).length))
tryGrow(es, cap);
try {
final Comparator<? super E> cmp;
if ((cmp = comparator) == null)
siftUpComparable(n, e, es);
else
siftUpUsingComparator(n, e, es, cmp);
size = n + 1;
notEmpty.signal();
} finally {
lock.unlock();
}
return true;
}

先取得主锁,查看当前 size 和数组长度;如果数组满,进入 tryGrow()。真正上浮时仍然持有主锁,因此不会让其他生产者、消费者同时修改同一份堆结构。

元素入堆、size 更新完成之后,notEmpty.signal() 通知等待消费者。队列可能已经有其他成员,这个通知也能让另一个等待线程重新获得机会。

它的 put() 直接调用 offer();超时 offer() 也直接调用 offer(),源码不使用 timeout 和 unit 来等待容量。这一点与 ArrayBlockingQueue、LinkedBlockingQueue 完全不同,虽然它们实现相同的 BlockingQueue 接口。

因此,选一个实现不能只看“有 put() 方法”。有界队列的 put() 表达等待空位;优先级无界队列的 put() 表达通常可以直接接收、无需容量等待。

tryGrow:为什么分配时先释放主锁

源码定位:PriorityBlockingQueue.java:287–310。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
// 先释放主锁尝试分配,再重新取得主锁并核对数组版本后发布。
private void tryGrow(Object[] array, int oldCap) {
lock.unlock(); // 必须先释放主锁,之后再重新获取
Object[] newArray = null;
if (allocationSpinLock == 0 &&
ALLOCATIONSPINLOCK.compareAndSet(this, 0, 1)) {
try {
int growth = (oldCap < 64)
? (oldCap + 2) // 小数组增长更快
: (oldCap >> 1);
int newCap = ArraysSupport.newLength(oldCap, 1, growth);
if (queue == array)
newArray = new Object[newCap];
} finally {
allocationSpinLock = 0;
}
}
if (newArray == null) // 其他线程正在分配时退让
Thread.yield();
// 先获取保护队列结构的锁。
lock.lock();
if (newArray != null && queue == array) {
queue = newArray;
System.arraycopy(array, 0, newArray, 0, oldCap);
}
}

数组分配可能比较慢。如果一直持有主锁分配,消费者即使本来有元素可取,也要等待扩容完成。源码因此先 unlock(),让其他线程可以继续操作原来的堆。

allocationSpinLock 用 CAS 抢占分配资格。拿到资格的线程计算新容量、分配新数组,最后恢复分配标记;没有分配到数组的线程可能 yield(),再回到主锁和外层容量检查。

重新获取主锁后,必须确认 queue 仍然是之前的 array。其他线程可能已经完成扩容,过时的新数组不能覆盖新的当前数组。

通过确认后,才替换 queue,并把旧数组内容复制到新数组。这里的旧数组在锁释放期间可能被合法的消费修改,最终复制在重新取得主锁之后进行,因此不是把释放锁前的旧内容快照盲目盖回去。

首选增长策略同样区分旧容量小于 64 与更大数组:前者增加 oldCap + 2,后者增加 oldCap >> 1,再交给 ArraysSupport 处理边界。

这个机制优化的是主锁占用和冗余分配竞争,不意味着整个 PriorityBlockingQueue 变成了无锁队列。真正的堆成员操作仍依赖主锁。

take:空堆等待,非空堆下沉修复

源码定位:PriorityBlockingQueue.java:529–540。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
// 没有成员时等待 notEmpty;dequeue 在锁内完成堆顶删除与下沉。
public E take() throws InterruptedException {
final ReentrantLock lock = this.lock;
// 获取锁时也响应线程中断。
lock.lockInterruptibly();
E result;
try {
while ( (result = dequeue()) == null)
// 暂时没有可取元素,释放锁等待生产者通知。
notEmpty.await();
} finally {
lock.unlock();
}
return result;
}

dequeue() 尝试取得 queue[0]。非空时保存堆顶,拿最后一个成员填根,清空旧尾槽位并下沉;为空则返回 null。

take() 用 while 包围 dequeue() 与 await()。条件被唤醒并不是某个元素已经专门分配给当前消费者,醒来后仍需要在锁内重新尝试取得。

普通 poll() 不等元素,锁内执行 dequeue() 后返回;超时 poll() 则按剩余预算等待 notEmpty。peek() 也取得主锁读取堆顶,保证与当前结构修改协调。

1
2
3
4
5
6
PriorityBlockingQueue<Integer> queue = new PriorityBlockingQueue<>();
System.out.println(queue.poll(10, TimeUnit.MILLISECONDS)); // null
queue.offer(2);
queue.offer(1);
System.out.println(queue.peek()); // 1,只查看
System.out.println(queue.poll()); // 1,已经删除

取出一个成员之后,其他更高优先级成员可能马上入队。单个方法按其操作时的堆状态完成,不承诺调用者处理返回值期间,队列里不会出现更优先的元素。

相同优先级的稳定次序

如果两个任务只按 priority 比较,比较结果为 0,队列仍然保留两份任务,但不保证它们按插入顺序出队。

可以增加业务顺序号:

1
2
3
4
5
6
record Job(int priority, long order, String name) {}
PriorityBlockingQueue<Job> queue = new PriorityBlockingQueue<>(11,
Comparator.comparingInt(Job::priority).thenComparingLong(Job::order));
queue.put(new Job(1, 2, "B"));
queue.put(new Job(1, 1, "A"));
System.out.println(queue.take().name()); // A

顺序号为比较规则的一部分,所以更新已有任务的 priority/order 仍然不能直接改字段。队列不会自动感知并重新排列,需要移除后重插或使用新的不可变成员。

而且“消费者先取得 A”也不保证“A 的业务处理一定先完成”。多个消费者的处理时间不同,出队顺序与完成顺序必须分别协调。

查找、删除与批量移出

contains()/remove(Object) 在主锁内沿数组按 equals() 搜索,不根据比较器寻找“同优先级任务”。找到后再执行堆修复,所以完整指定删除通常为 O(n)。

drainTo() 在锁内反复 dequeue(),因此移出的成员按当时堆的优先级逐个选择。这与 iterator() 沿数组访问不同;转移 k 个成员时,堆调整会产生约 O(k log n) 的成本。

1
2
3
4
5
6
PriorityBlockingQueue<Integer> queue = new PriorityBlockingQueue<>();
queue.addAll(List.of(30, 10, 20));
List<Integer> batch = new ArrayList<>();
queue.drainTo(batch, 2);
System.out.println(batch); // [10, 20]
System.out.println(queue.peek()); // 30

drainTo() 的目标 add() 失败仍然不提供两个容器之间的自动全部回滚;整个过程还会在队列锁内调用目标添加操作,实际目标的成本会影响主锁占用。

size() 返回当前成员数,remainingCapacity() 固定返回 Integer.MAX_VALUE。不能通过后者计算还可以积压多少字节,也不能用 size() 先判断再 offer() 实现精确并发上限。

如果业务需要同时具备“优先级”和“严格限制待处理数量”,应在队列之外建立一致的准入协议,或者选择专门的有界优先级实现。准入令牌还必须与删除、取消和消费的生命周期配套,不能只在成功添加时减一、却在某些退出路径忘记释放。

序列化使用临时 PriorityQueue 保存当前逻辑成员与比较规则,恢复时重新建立内部数组堆。这个辅助字段的存在服务于序列化,不改变正常读写的主锁加数组结构。

迭代器

源码的 iterator() 返回 new Itr(toArray())。toArray() 在锁下复制当前有效数组成员,所以这个具体实现的迭代器保存的是创建时的成员引用快照。

创建迭代器为 O(n) 复制,之后普通推进为 O(1),完整遍历 O(n)。这与 CopyOnWriteArraySet 创建时只保存已经存在的数组引用不同,不能把两种快照的创建成本都写成 O(1)。

接口与文档把这种遍历归入弱一致支持范围,但是源码给出的实际观察行为更具体:创建后加入的新成员不进入该数组,后来删除的成员仍可能被快照返回。

1
2
3
4
5
6
7
8
9
10
PriorityBlockingQueue<Integer> queue = new PriorityBlockingQueue<>();
queue.addAll(List.of(3, 1, 2));
Iterator<Integer> old = queue.iterator();
queue.clear();
queue.offer(9);
List<Integer> snapshot = new ArrayList<>();
old.forEachRemaining(snapshot::add);
Collections.sort(snapshot);
System.out.println(snapshot); // [1, 2, 3]
System.out.println(queue.peek()); // 9

快照数组来自堆布局,原始遍历不保证排序。这里主动 sort() 只是为了展示快照成员,不应拿它作为迭代器有序的证据。

它支持 iterator.remove():使用快照中最后返回的对象,通过 removeEq() 按引用身份删除当前堆中的对应成员,再重置 lastRet。

源码定位:PriorityBlockingQueue.java:664–678。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
// 按引用身份寻找快照所指的对象,删除时仍需恢复当前堆性质。
void removeEq(Object o) {
final ReentrantLock lock = this.lock;
// 先获取保护队列结构的锁。
lock.lock();
try {
final Object[] es = queue;
for (int i = 0, n = size; i < n; i++) {
if (o == es[i]) {
removeAt(i);
break;
}
}
} finally {
lock.unlock();
}
}

这里的 o == array[i] 与公开 remove(Object) 的 equals() 匹配不同。两个对象逻辑相等,但快照返回的是其中一个实例,迭代器就不应该把另一个实例当作同一观察对象删除。

如果同一个对象引用多次加入队列,身份比较仍只能定位一个当前匹配槽位,并不建立一套不可复用的“入队序号”。对象引用、对象相等和一次具体提交的身份,业务上如果需要区分,应另外保存唯一标识。

源码定位:PriorityBlockingQueue.java:840–842。

1
2
3
4
// 创建时复制有效数组成员,后续推进读取这份引用副本。
public Iterator<E> iterator() {
return new Itr(toArray());
}

源码定位:PriorityBlockingQueue.java:860–864。

1
2
3
4
5
6
// cursor 沿快照数组移动,因此后续堆重排不改变这次遍历的槽位。
public E next() {
if (cursor >= array.length)
throw new NoSuchElementException();
return (E)array[lastRet = cursor++];
}

next() 中没有获取主锁,因为它读的是已经复制完成的数组,不是正在被 offer()/poll() 修改的原数组。只有 iterator.remove() 会重新进入当前集合的同步路径;“迭代器无需逐步持锁”与“创建快照时需要持锁”要分别看。

如果快照成员早已离开当前队列,就不会再凭 equals() 任意删除一份逻辑相等的其他对象。它也不会把删除动作限制为只修改快照数组,删除影响的是当前集合。

快照隔离的是成员引用排列,不是业务对象字段。对象本身可变时,快照仍然可以看到同一个对象后续的变化,集合不负责这些字段的同步。

Spliterator 在绑定时也会通过数组副本工作,报告 SIZED、SUBSIZED、NONNULL,不报告 ORDERED/SORTED。支持并行遍历不等于业务任务会按优先级并行完成,也不等于容器的原始成员数组已经有全局排序。

DelayQueue

基本特性

PriorityQueue 决定的是“谁更优先”,DelayQueue 又增加了一个条件:排在最前面的元素,只有到期之后才能通过通常的出队方法取得。

它继承 AbstractQueue,实现 BlockingQueue<E>,其中 E 必须实现 Delayed。它允许重复,禁止 null,线程安全,底层使用 PriorityQueue,并以一把 ReentrantLock 协调访问。

DelayQueue 是逻辑无界队列,put() 不等待容量,超时 offer() 也不会为容量使用传入的等待参数。remainingCapacity() 返回 Integer.MAX_VALUE,真实存储仍会受内存资源限制。

它并不是加入任务之后自动启动线程执行任务。DelayQueue 只保存和交付到期元素,何时调用 take()、取得之后如何处理,都由使用它的代码负责。

常见用途是延后执行、缓存过期消息、超时检查。它只提供队列层面的时间门槛,不提供精确到某个时刻必须执行完毕的实时保证。

“队列非空”与“有可以出队的元素”也不一样:size() 可以大于 0,peek() 可以返回未来任务,但是 poll() 仍然返回 null,take() 仍然需要等待。

结构分析

主要字段为 PriorityQueue<E> q、ReentrantLock lock、Condition available,以及 Thread leader。

q 负责把最早应该到期的元素放到堆顶;lock 保护堆和等待协调状态;available 用于等待“有元素”或“当前堆顶已经可以取得”。

leader 不是额外创建的常驻调度线程,而是当前负责按堆顶剩余延迟定时等待的消费者线程。

DelayQueue 的堆顶时间门槛与 leader 等待

为什么需要 leader?如果有十个消费者,而最早任务将在五秒后到期,让十个线程都设置五秒等待,会产生冗余定时唤醒与竞争。

源码让一个线程作为 leader 定时等待,其他消费者等待 available 通知。leader 退出或新堆顶出现时,再安排后续线程重新检查并接任。

这并没有保证永远只有一个线程处于所有形式的等待。其他线程可以无限期等待条件,带超时的 poll() 也有调用者自己的等待预算;核心是减少针对同一个堆顶到期时刻的重复定时等待。

基本操作

Delayed:排序与剩余延迟必须表达同一个时间关系

Delayed 继承 Comparable<Delayed>,提供 getDelay(TimeUnit)。compareTo() 决定堆的位置,getDelay() 决定堆顶是否已经到期。

这两个方法必须保持一致的含义。如果 compareTo() 把一个很晚到期的任务排在前面,后面的已到期任务就可能被堆顶阻挡,因为出队只检查堆顶,而不会扫描所有成员寻找另一份已到期任务。

通常用固定的截止时间实现,而不是在每次 compareTo() 时依赖不稳定的业务字段。计算相对等待时,System.nanoTime() 比 wall clock 更合适,因为它用于测量经过时间,不受系统时钟校准直接改变。

下面是一个用于短时延迟的完整元素示例:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
class Ticket implements Delayed {
final String name;
final long deadline;

Ticket(String name, long delay, TimeUnit unit) {
long nanos = unit.toNanos(delay);
if (delay < 0 || nanos > Long.MAX_VALUE / 4)
throw new IllegalArgumentException("delay out of example range");
this.name = name;
this.deadline = System.nanoTime() + nanos;
}

@Override
public long getDelay(TimeUnit unit) {
// 返回截止时间相对于当前时刻的剩余量。
return unit.convert(deadline - System.nanoTime(), TimeUnit.NANOSECONDS);
}

@Override
public int compareTo(Delayed other) {
// 示例中的时间跨度受限,用差值比较避免简单转成 int。
long diff = deadline - ((Ticket) other).deadline;
return Long.compare(diff, 0L);
}
}

DelayQueue<Ticket> queue = new DelayQueue<>();
queue.offer(new Ticket("later", 1, TimeUnit.DAYS));
System.out.println(queue.peek().name); // later,未到期也能查看
System.out.println(queue.poll()); // null,不能普通出队
System.out.println(queue.size()); // 1

这个示例约束了短时间跨度,比较差值的前提是参与比较的截止时间差不会跨越 long 的半范围。不能将任意超大延迟、长期陈旧任务的数学关系都无条件塞进 long 加减。

实际业务若涉及绝对日期,还需要明确怎样将日期转换为等待期限,怎样处理重启与时钟变化。nanoTime() 的起点没有业务日期含义,不能把它直接当作时间戳持久化后在另一进程继续使用。

添加元素:新堆顶要使旧 leader 计划失效

源码定位:DelayQueue.java:166–179。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
// 在锁内入堆;新元素成为堆顶时,解除旧 leader 标记并通知一个等待者。
public boolean offer(E e) {
final ReentrantLock lock = this.lock;
// 先获取保护队列结构的锁。
lock.lock();
try {
q.offer(e);
if (q.peek() == e) {
leader = null;
available.signal();
}
return true;
} finally {
lock.unlock();
}
}

如果 e 成为 q.peek(),说明此前围绕旧堆顶建立的定时等待计划可能不再合适。源码将 leader 设为 null,并对 available 调用 signal(),让某个线程重新检查最早截止时间。

这里比较的是对象身份 q.peek() == e,不是 equals()。它判断刚加入的这份对象是否成为堆顶,不是寻找任意相等任务。

比如原堆顶一分钟后到期,新加入的任务现在就能出队,就不能等旧 leader 的一分钟预算走完才处理。通知与 leader 重置使等待安排可以及时调整。

put() 和带超时的 offer() 都委派给普通 offer()。它们不等待空位,但入堆仍需锁,也可能触发底层 PriorityQueue 扩容。

poll、peek 与异常式出队

源码定位:DelayQueue.java:214–225。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
// 只有堆顶存在且剩余延迟不大于 0,才真正从 PriorityQueue 取出。
public E poll() {
final ReentrantLock lock = this.lock;
// 先获取保护队列结构的锁。
lock.lock();
try {
E first = q.peek();
return (first == null || first.getDelay(NANOSECONDS) > 0)
? null
: q.poll();
} finally {
lock.unlock();
}
}

peek() 则只查看 q.peek(),不检查是否到期。所以 peek() 的非 null 返回值,只能说明当前存在一个堆顶,不能直接证明 poll() 一定成功。

无参数 remove() 在没有已到期堆顶时会抛 NoSuchElementException,即使 size() 大于 0。element() 作为查看方法,则沿 peek() 的含义查看当前堆顶。

没有参数的 remove(),与 remove(Object) 的限制不同。 前者取得到期队首,后者用于主动删除指定任务,后者不要求任务已经到期。

take:leader 与 follower 怎样协作

源码定位:DelayQueue.java:235–267。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
// 循环检查堆顶,只有一个当前 leader 围绕剩余延迟进行定时条件等待。
public E take() throws InterruptedException {
final ReentrantLock lock = this.lock;
// 获取锁时也响应线程中断。
lock.lockInterruptibly();
try {
for (;;) {
E first = q.peek();
if (first == null)
available.await();
else {
long delay = first.getDelay(NANOSECONDS);
if (delay <= 0L)
return q.poll();
first = null; // 等待期间不要保留该引用
if (leader != null)
available.await();
else {
Thread thisThread = Thread.currentThread();
leader = thisThread;
try {
// 只有 leader 按当前堆顶的剩余延迟定时等待。
available.awaitNanos(delay);
} finally {
if (leader == thisThread)
leader = null;
}
}
}
}
} finally {
if (leader == null && q.peek() != null)
available.signal();
lock.unlock();
}
}

把这个循环分开看就清楚了。

队列为空时,等待 available;堆顶已经到期时,直接 q.poll() 返回;堆顶还没到期且已有 leader 时,当前线程作为 follower 等通知;没有 leader 时,当前线程成为 leader,等待当前 delay。

available.await()/awaitNanos() 会释放 lock,否则生产者无法添加更早任务,其他消费者也无法协作。醒来后重新取得锁,再次检查 q.peek() 与 delay。

源码还有 first = null,避免等待期间在局部变量中额外长期保留旧堆顶引用。队列自身需要保存未来任务,但等待线程没有必要继续强持有一个可能已经被移除的旧任务。

leader 的 finally 只在 leader == thisThread 时清空标记。因为等待期间其他线程可能因新堆顶重置或接管了 leader,旧线程不能随意清掉别人的当前状态。

最外层 finally 则在没有 leader 且 q 仍非空时,通知一个等待者继续组织下一轮定时等待,再释放锁。

这套协议把堆顶变化、定时等待和通知联系起来,不是“线程每隔固定时间醒来 poll() 一下”。它能够避免大量无意义的周期扫描。

带超时的 poll:同时考虑调用预算和任务期限

poll(timeout, unit) 不仅要看堆顶 delay,还要看调用者剩余 nanos。如果队列为空,就按自身预算等待;堆顶比预算更晚到期,或者已有 leader 时,也可能只按预算等待。

如果当前线程有条件接任 leader,而且堆顶会在预算内到期,它就围绕堆顶 delay 等待,并据实际耗时扣减自身预算。

因此,不能把这个方法简化成“take() 然后最多 sleep() timeout”。它需要在每一轮唤醒后同时检查成员、到期状态、leader 状态与剩余时间。

timeout 为 0 也不意味着忽略已经到期的成员。方法先检查当前状态,若有已到期堆顶,可以直接取得;只有需要继续等条件时,才由耗尽预算决定返回 null。

源码定位:DelayQueue.java:280–319。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
// 在同一循环中协调调用预算、堆顶延迟和 leader 身份。
public E poll(long timeout, TimeUnit unit) throws InterruptedException {
long nanos = unit.toNanos(timeout);
final ReentrantLock lock = this.lock;
// 获取锁时也响应线程中断。
lock.lockInterruptibly();
try {
for (;;) {
E first = q.peek();
if (first == null) {
if (nanos <= 0L)
return null;
else
nanos = available.awaitNanos(nanos);
} else {
long delay = first.getDelay(NANOSECONDS);
if (delay <= 0L)
return q.poll();
if (nanos <= 0L)
return null;
first = null; // 等待期间不要保留该引用
if (nanos < delay || leader != null)
nanos = available.awaitNanos(nanos);
else {
Thread thisThread = Thread.currentThread();
leader = thisThread;
try {
// 只有 leader 按当前堆顶的剩余延迟定时等待。
long timeLeft = available.awaitNanos(delay);
nanos -= delay - timeLeft;
} finally {
if (leader == thisThread)
leader = null;
}
}
}
}
} finally {
if (leader == null && q.peek() != null)
available.signal();
lock.unlock();
}
}

这里的 nanos -= delay - timeLeft 不是简单把原来的 delay 全部扣掉,因为 awaitNanos() 可能提前被通知唤醒,返回的 timeLeft 反映剩余等待量。源码按这一轮实际消耗更新总预算,再回到循环读取可能已经变化的堆顶。

线程成为 leader 并不意味着旧堆顶永久属于它。等待时锁已经释放,其他线程可以取消旧任务、加入更早任务或取得已到期成员;只有重新获得锁并再次检查,才能决定下一步。

主动取消与批量移出

remove(Object)、iterator.remove() 和 clear() 不要求成员已经到期,可用于取消任务或清理整个容器。主动取消是按对象匹配删除,不是把该任务的 delay 设成 0。

drainTo() 则只移出当前已到期的队首,遇到未到期堆顶就停止。它不会因为“批量操作”就绕过时间门槛,把全部未来任务也交出去。

1
2
3
4
DelayQueue<Delayed> queue = new DelayQueue<>();
// 实际使用时加入自己的 Delayed 实现。
List<Delayed> expired = new ArrayList<>();
System.out.println(queue.drainTo(expired)); // 空队列:0

clear() 不等待队列元素自然到期,直接清空底层堆。对正在等待的消费者而言,“当前没有成员”仍可能意味着继续等待未来添加,不能把 clear() 当作统一关闭信号。

延迟字段如果在入堆后直接修改,堆也不会自动重排。某任务被改成更早到期,但仍留在旧位置,可能继续被旧堆顶阻挡;需要先移除,更新或创建新对象,再重新加入。

迭代器

DelayQueue 的 iterator() 返回 new Itr(toArray())。toArray() 在锁保护下复制当前堆成员,所以这份具体实现保存创建时的成员引用数组。

迭代器包含已到期和未到期成员,不受普通出队的时间门槛限制,也不保证按到期时间排序,因为快照保留的是堆数组布局。

1
2
3
普通 poll/take:检查当前堆顶是否到期,再移除
iterator:访问创建时复制的所有成员,不检查 getDelay
iterator.remove:主动取消当前对应对象,不要求到期

它没有 expectedModCount,不因集合后续增删就抛 Fail-Fast 异常。创建之后的新任务不进入此快照,已经从当前队列删除的任务仍可能被旧快照返回。

iterator.remove() 调用 removeEQ(),取得 lock 后扫描底层堆,按对象身份找到快照里最后返回的任务,再使用底层迭代删除机制。

这样可以避免 equals() 相等但属于另一任务实例的成员被误删。删除影响当前队列,而不是只把迭代器数组中的位置变成 null;一次 next() 最多对应一次 remove()。

创建快照需要 O(n) 复制,普通遍历 O(n),一次按身份取消还需 O(n) 扫描加堆修复。快照稳定的仍是成员引用,业务对象自己的字段没有深拷贝,也没有自动同步。

它的类并没有直接声明 Serializable,因此不能仅凭底层 PriorityQueue 可序列化,就认定整个 DelayQueue 也具有相同的序列化接口能力。任务跨重启恢复,需要单独设计逻辑期限与持久化方案。

SynchronousQueue

基本特性

如果说前面的阻塞队列是在生产者和消费者之间放了一个“仓库”,那么 SynchronousQueue 就更像是一次当面交接:生产者把元素交给消费者,元素不会作为普通成员保存在队列中,等待下一次取出。

它实现 BlockingQueue,不接受 null,支持多线程并发访问。最特别的地方是 没有用于积压元素的容量,而不是“容量只有 1”。如果没有能够立即匹配的消费者,普通 offer() 会返回 false;如果没有能够立即匹配的生产者,普通 poll() 会返回 null。

put() 与 take() 则允许等待。put() 返回通常表示这次交付已经完成匹配,而不是只把元素放进一个缓冲区。不过,匹配完成也不代表消费者的业务逻辑已经执行完毕:消费者取得任务后,还可能需要很长时间处理。

1
2
3
4
5
SynchronousQueue<String> queue = new SynchronousQueue<>();
System.out.println(queue.offer("任务")); // false:没有正在接收的线程
System.out.println(queue.poll()); // null:没有正在交付的线程
System.out.println(queue.size()); // 0
System.out.println(queue.remainingCapacity()); // 0

这里 offer() 返回 false,元素也没有留下来,因此稍后的 poll() 不会拿到刚才那个任务。不能把失败的 offer() 当成“暂时还没处理,但已经提交成功”。

构造时可以选择公平策略。默认构造器使用非公平策略;传入 true 时,等待交接的线程按 FIFO 次序获得匹配机会。这个顺序针对已经进入等待结构的参与者,不能扩大为“所有线程从调用方法开始就严格按照现实时间先后返回”。

SynchronousQueue 经常用于希望生产与消费直接配对的场景,也适合不希望队列形成任务积压的执行器。它能够把“接收方尚未准备好”直接反馈给提交方,但具体反馈是等待、超时还是立即失败,仍然取决于调用哪个方法。

结构分析

这里需要特别注意源码版本。本文核对的是本机 JDK 21.0.4 的 src.zip 及运行类签名:这个版本使用 Transferer 继承 LinkedTransferQueue,并在其上增加 LIFO 模式支持。不能拿旧版本中 TransferStack、TransferQueue 两套独立内部类的结构,直接替代当前代码。

SynchronousQueue 自身保留 transferer 和 fair 字段,操作再根据公平策略分派:

源码定位:SynchronousQueue.java:231–234。

1
2
3
4
5
// 公平模式走 FIFO 双重队列,默认模式走 LIFO 等待结构。
private Object xfer(Object e, long nanos) {
Transferer<E> x = transferer;
return (fair) ? x.xfer(e, nanos) : x.xferLifo(e, nanos);
}

e 非 null 表示生产者带着数据到达;e 为 null 表示消费者请求数据。这个约定也解释了为什么用户不能插入 null:否则内部无法区分“交付 null”和“请求一个元素”。

nanos 为 0 表示只尝试立即匹配;Long.MAX_VALUE 表示没有指定超时期限;一个有限的正值表示限时等待。SynchronousQueue 不使用 LinkedTransferQueue 的异步入队模式,因为它不允许将未交接元素作为普通缓冲成员留下。

SynchronousQueue 的两种交接顺序

公平模式把同一类型的等待者链接在队尾,再由互补类型的到达者匹配较早的等待节点。默认模式则从等待结构顶端寻找互补节点,同类参与者需要等待时压到顶端,因此匹配顺序表现为 LIFO。

图中“数据节点”和“请求节点”是协调线程的内部记录,不能等同于普通集合成员。一个等待生产者的节点虽然暂时持有任务引用,对外 size()、peek()、iterator() 依然表现为空。

下面是当前版本的 LIFO 核心。它不仅处理一次正常匹配,还要跳过已经匹配的节点,处理 CAS 失败后的重试,以及取消等待后的清理:

源码定位:SynchronousQueue.java:167–201。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
// 先尝试匹配栈顶的互补节点;没有匹配者时才登记等待。
final Object xferLifo(Object e, long ns) {
boolean haveData = (e != null);
Object m; // 匹配结果;没有匹配时就是 e
outer: for (DualNode s = null, p = head;;) {
while (p != null) {
boolean isData; DualNode n, u; // 协助收缩
if ((isData = p.isData) != ((m = p.item) != null))
p = (p == (u = cmpExHead(p, (n = p.next)))) ? n : u;
else if (isData == haveData) // 同类模式,压到下面
break;
else if (p.cmpExItem(m, e) != m)
p = head; // 错过匹配,重新开始
else { // 匹配到互补节点
Thread w = p.waiter;
cmpExHead(p, p.next);
// 通知等待线程继续检查匹配状态。
LockSupport.unpark(w);
break outer;
}
}
if (ns == 0L) { // 没有匹配,也不等待
m = e;
break;
}
if (s == null) // 尝试压入节点并等待
s = new DualNode(e, haveData);
s.next = p;
if (p == (p = cmpExHead(p, s))) {
if ((m = s.await(e, ns, this, // 队列几乎为空时先自旋
p == null || p.waiter == null)) == e)
unspliceLifo(s); // 已取消
break;
}
}
return m;
}

其中 p.cmpExItem(m, e) != m 表示原子比较交换失败:节点状态已经被其他线程改变,当前线程必须重新寻找匹配对象。成功时,数据节点的 item 从数据变为 null,或者请求节点的 item 从 null 变为数据,交接在这次原子状态转换中完成。

unpark() 用于唤醒已经挂起的等待者,但唤醒本身不是匹配的依据。等待线程必须再次检查节点状态,判断究竟是正常交接、超时取消还是中断取消。

取消节点也不能简单地随意改写 next。另一个线程可能正在用相同的前驱和后继执行匹配或清理,源码通过原子更新、跳过失效节点及周期清扫共同推进结构。某些已经失效的节点可以暂时留在链上,而不会重新成为有效数据。

基本操作

1. put:等到元素完成交接

源码定位:SynchronousQueue.java:261–269。

1
2
3
4
5
6
7
8
9
10
// 无限期交接使用 Long.MAX_VALUE,未成功匹配且被中断时抛出异常。
public void put(E e) throws InterruptedException {
Objects.requireNonNull(e);
if (!Thread.interrupted()) {
if (xfer(e, Long.MAX_VALUE) == null)
return;
Thread.interrupted(); // 只可能因中断而失败
}
throw new InterruptedException();
}

put() 不执行“检查当前数量,再放入数组”的逻辑,而是把本次数据提交给交接算法。如果已有等待消费者,可以立即匹配;否则登记数据节点,在 await() 中等待。

这里所谓无限期,指没有由调用者提供超时期限,并不意味着线程只能永远挂起。正常匹配或中断都会使它离开等待。调用方应该保留中断传播或恢复中断状态,避免把程序的取消信号吞掉。

可以用一个确定存在消费者的例子观察交接:

1
2
3
4
5
6
7
8
9
10
SynchronousQueue<String> queue = new SynchronousQueue<>();
ExecutorService executor = Executors.newSingleThreadExecutor();
try {
Future<String> received = executor.submit(queue::take);
queue.put("订单-A");
System.out.println(received.get()); // 订单-A
System.out.println(queue.size()); // 0
} finally {
executor.shutdownNow();
}

无论消费者先开始等待还是生产者先到达,两个阻塞操作都能在互补操作到达时配对。例子中的 Future.get() 是为了取得消费者线程的结果;它不是 SynchronousQueue 自动等待业务处理完成的功能。

2. take:等到获得一次交付

源码定位:SynchronousQueue.java:313–321。

1
2
3
4
5
6
7
8
9
10
// 消费者用 null 表示请求;获得非 null 数据说明匹配成功。
public E take() throws InterruptedException {
Object e;
if (!Thread.interrupted()) {
if ((e = xfer(null, Long.MAX_VALUE)) != null)
return (E) e;
Thread.interrupted();
}
throw new InterruptedException();
}

如果队列中已有等待的生产者,take() 会获得其中一个交付的数据。否则,消费者登记请求节点并等待生产者到来。因此,一个空队列上的 take() 会等待,却不代表队列内部什么节点都没有。

生产者与消费者不需要持有同一把全局排他锁完成整个操作。当前实现依赖节点状态的原子更新协调交接;需要等待时,还会使用短暂自旋和 LockSupport 等机制。

这也说明 基于 CAS 的实现可以提供阻塞方法。CAS 描述结构更新方式,put()/take() 的阻塞描述 API 等待条件,两者并不矛盾。await() 还针对虚拟线程避免沿用平台线程的自旋策略,并在适用时通过 ForkJoinPool.ManagedBlocker 协调工作线程阻塞。

核心等待方法复用自 LinkedTransferQueue.DualNode:

源码定位:LinkedTransferQueue.java:423–465。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
// 在短暂自旋、登记 waiter、挂起与取消之间协调,始终以节点状态判断结果。
final Object await(Object e, long ns, Object blocker, boolean spin) {
Object m; // 匹配结果;没有匹配时就是 e
boolean timed = (ns != Long.MAX_VALUE);
long deadline = (timed) ? System.nanoTime() + ns : 0L;
boolean upc = isUniprocessor; // 不自旋,但稍后重新检查
Thread w = Thread.currentThread();
if (w.isVirtual()) // 不自旋
spin = false;
int spins = (spin & !upc) ? SPINS : 0; // 为负时可以挂起
while ((m = item) == e) {
if (spins >= 0) {
if (--spins >= 0)
Thread.onSpinWait();
else { // 准备挂起
if (spin) // 偶尔重新检查
checkForUniprocessor(upc);
LockSupport.setCurrentBlocker(blocker);
waiter = w; // 保证顺序
VarHandle.fullFence();
}
} else if (w.isInterrupted() ||
(timed && // 用不可能的匹配尝试取消
((ns = deadline - System.nanoTime()) <= 0L))) {
m = cmpExItem(e, (e == null) ? this : null);
break;
} else if (timed) {
if (ns < SPIN_FOR_TIMEOUT_THRESHOLD)
Thread.onSpinWait();
else
// 暂时挂起当前线程,醒来后仍需检查状态。
LockSupport.parkNanos(ns);
} else if (w instanceof ForkJoinWorkerThread) {
try {
ForkJoinPool.managedBlock(this);
} catch (InterruptedException cannotHappen) { }
} else
LockSupport.park();
}
if (spins < 0) {
LockSupport.setCurrentBlocker(null);
waiter = null;
}
return m;
}

等待前登记 waiter,让成功匹配的另一方能够找到需要唤醒的线程。登记之后还必须再次读 item,以处理“匹配发生在登记前后”的竞速,避免已经成功交付却把自己挂起到没有后续通知的状态。

中断或期限耗尽时,线程通过比较交换尝试把节点从未匹配状态改为取消状态。这个操作可能失败,因为对方已经先完成匹配;源码最终返回的状态决定调用结果,而不是由两个线程各自推测是否已经超时。

源码中的 park()/parkNanos() 允许无理由返回,所以等待循环同样不能只挂起一次就相信交付成功。清理 waiter 与阻塞器记录,也用于避免节点长时间保留已经不再等候的线程引用。

3. offer/poll:立即尝试与限时交接

源码定位:SynchronousQueue.java:300–303。

1
2
3
4
5
// 立即交接没有匹配者便失败,不把元素留作队列成员。
public boolean offer(E e) {
Objects.requireNonNull(e);
return xfer(e, 0L) == null;
}

普通 offer() 不等待将来的消费者。即使另一个消费者“一会儿就会来”,本次调用仍然可能失败。普通 poll() 同理,不能因为应用中存在生产线程,就假定此刻一定能取得任务。

1
2
3
4
SynchronousQueue<Integer> queue = new SynchronousQueue<>();
boolean accepted = queue.offer(10, 10, TimeUnit.MILLISECONDS);
System.out.println(accepted); // false:没有消费者
System.out.println(queue.poll(10, TimeUnit.MILLISECONDS)); // null

带超时的方法允许节点在等待结构中保留到匹配成功、超时或中断。超时取消与对方匹配可能竞速,源码由原子节点状态决定谁先成功,不能仅靠墙上时钟判断方法必然返回哪一种结果。

超时值是等待策略,也不是精准调度保证。线程还要受操作系统调度影响,不能用超时 API 实现严格的实时控制。

对于普通 offer()/poll(),“立即”表示不为未来交接等待条件;当前线程仍可能因线程调度、原子操作竞争和重试花费时间,不能理解成固定的几条指令或绝对不会延迟。

4. 集合查询为何全部看起来为空

size() 始终返回 0,isEmpty() 始终返回 true,remainingCapacity() 返回 0,peek() 返回 null。contains() 不会把等待生产者的数据当成可查询成员,remove(Object) 也不能按普通队列思路撤销另一个线程的交接。

clear() 对这种始终为空的集合视图没有普通“清空任务仓库”的意义,它不能被当作终止所有等待线程的接口。若要取消正在 put()/take() 的线程,应使用中断或应用定义的生命周期方案。

这里还有一个容易误判的操作:drainTo() 通过反复 poll() 来完成交接。如果此时有生产者已经等待,poll() 可以取得其数据。因此,drainTo() 可能转移元素,即使调用前的 size() 为 0。

它转移的是当下可以立即匹配的交付,而不是读取一个隐藏的缓冲区。maxElements 限制本轮最多接收多少个元素,也不建立固定大小的储存空间。

迭代器

源码定位:SynchronousQueue.java:462–464。

1
2
3
4
// 对外集合视图始终为空,因此返回空迭代器。
public Iterator<E> iterator() {
return Collections.emptyIterator();
}

SynchronousQueue 的迭代器没有待遍历成员。hasNext() 返回 false,next() 抛出 NoSuchElementException;它不会等待,也不会把内部等待节点暴露成元素。

因此,不能用 for-each 监控还在等待交接的任务,更不能用 iterator.remove() 解除生产者阻塞。要观察提交是否成功,应看 offer() 的返回值、put() 是否返回、超时结果或者应用额外记录的状态。

spliterator() 同样为空。toArray() 的结果也是空数组,对外的集合语义与 size()、peek() 保持一致。这里不宜套用“并发队列的迭代器都是弱一致”的结论:本类的重点是 空集合视图。

SynchronousQueue 实现 Serializable,但序列化不会保存正在交接中的业务线程或把它们迁移到另一个进程。源码中为序列化兼容保留的锁和等待队列字段,也不能被误当作当前正常交接路径使用的全局锁。

回到使用选择:希望积压一定数量的任务时,可以考虑有容量的阻塞队列;希望提交方与接收方直接相遇时,SynchronousQueue 才能表达这种约束。它的关键不在“队列长度总是 0”这个现象,而在每个成功提交背后都有一次已经完成的交接匹配。

ConcurrentLinkedQueue

基本特性

ConcurrentLinkedQueue 是一个基于链表的并发 FIFO 队列,实现 Queue,允许重复元素,不允许 null,没有构造时指定的容量上限。它没有 put()/take() 这类等待队列条件的接口:空队列上的 poll() 返回 null,offer() 正常情况下将元素链接进队列并返回 true。

它与 LinkedBlockingQueue 的区别不能只概括成“链表对链表”。LinkedBlockingQueue 用锁和条件变量提供容量限制及阻塞等待;ConcurrentLinkedQueue 的主要入队、出队路径通过 CAS 更新节点,不用一把排他锁保护整条链。

这里的非阻塞也不是“每个线程一定在固定时间内完成”。某个线程可能反复遭遇 CAS 失败而重试,算法讨论的推进保证和单个线程的公平性是不同问题。本类也不提供构造器上的公平策略选项。

1
2
3
4
5
6
7
Queue<String> queue = new ConcurrentLinkedQueue<>();
queue.offer("A");
queue.offer("A"); // 允许相等的重复元素
queue.offer("B");
System.out.println(queue.poll()); // A
System.out.println(queue.poll()); // A
System.out.println(queue.peek()); // B,查看不删除

FIFO 指有效成员的队列顺序,不能据此推导两个竞争线程从“开始调用 offer()”的时间先后必然相同。并发提交的顺序由真正成功链接节点的原子操作确定。

入队之前的操作与另一线程随后访问或移除该元素的操作之间具有文档规定的 happens-before 关系。因此,先初始化任务,再通过队列发布,可以让消费者看见初始化状态;入队之后继续无同步地改写同一个任务对象,则不是队列帮你解决的事情。

结构分析

核心是 Node 中的 item 与 next,以及队列的 head、tail。item 和 next 为 volatile,head、tail 也支持相应的原子更新。节点保存用户元素引用,不会复制元素对象。

ConcurrentLinkedQueue 的逻辑删除与指针推进

head 不一定直接指向第一个有值节点,tail 也不一定随时指向物理上的最后一个节点。源码允许这两个入口适度滞后,遍历时再寻找真正可用的位置,这样可以减少每次操作必须执行的竞争更新。

需要区分“节点还在链上”和“元素还在队列里”。当 item 已经被置为 null,这个节点就不再代表有效成员,即使它仍被前驱的 next 引用着。

删除采用两个层次:先完成逻辑删除,让 item 从用户元素变为 null;之后再由当前操作或者后续遍历尝试跳过失效节点。物理清理影响可达链的形状,成员是否存在则由 item 状态决定。

另一个特殊状态是节点 next 指向自己。它表示该节点已经从可用的头部路径脱离。若某个较慢线程仍持有旧节点,就通过这个标记识别自己落在旧链上,并从当前 head 或 tail 重新寻找路径。

源码定位:ConcurrentLinkedQueue.java:290–294。

1
2
3
4
5
6
7
// 头指针推进成功后让旧头自链接,帮助旧遍历识别失效路径。
final void updateHead(Node<E> h, Node<E> p) {
// assert h != null && p != null && (h == p || h.item == null);
// 字段仍等于预期值时才完成原子更新。
if (h != p && HEAD.compareAndSet(this, h, p))
NEXT.setRelease(h, h);
}

自链接还有内存可达性上的作用:旧节点不再一直通过 next 引用后面整条仍然活跃的链。否则,一个被慢线程或迭代器暂时保留的旧节点,可能间接保留大量本来不需要保留的后续对象。

但这里不是手动释放内存。Java 的垃圾回收仍然根据对象可达性回收节点;源码只是避免不必要的引用链,并为遍历恢复提供明确标记。

基本操作

1. offer:找到尾端并链接新节点

源码定位:ConcurrentLinkedQueue.java:354–381。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
// 真正入队发生在 next 从 null 原子改为新节点的那一步。
public boolean offer(E e) {
final Node<E> newNode = new Node<E>(Objects.requireNonNull(e));

for (Node<E> t = tail, p = t;;) {
Node<E> q = p.next;
if (q == null) {
// p 是尾节点
// 字段仍等于预期值时才完成原子更新。
if (NEXT.compareAndSet(p, null, newNode)) {
// 这次成功的 CAS 就是线性化点:
// e 由此成为队列中的元素,newNode 也由此“生效”。
if (p != t) // 每次跳两个节点;失败也没关系
TAIL.weakCompareAndSet(this, t, newNode);
return true;
}
// CAS 竞争输给其他线程,重新读取 next
}
else if (p == q)
// 我们已经脱离链表。如果 tail 未变,它同样已脱离,
// 此时需要跳回 head——从 head 出发总能到达所有存活节点;
// 否则新的 tail 是更合适的选择。
p = (t != (t = tail)) ? t : head;
else
// 每跳两下检查一次 tail 是否已更新。
p = (p != t && t != (t = tail)) ? t : q;
}
}

首先创建新节点并拒绝 null。随后从观察到的 tail 出发,检查 p.next。如果 next 为 null,p 就是当前候选尾节点,当前线程尝试用 CAS 把新节点链接到它后面。

只有成功修改这个 next,元素才真正进入队列。之后更新 tail 属于帮助后续操作更快定位尾端,即使它没有成功,已经链接的元素也不会因此消失。

如果 next 非 null,说明已有节点排在 p 后面,当前线程需要继续向后走,或利用更新后的 tail 跳转。如果 p.next 等于 p,说明落入旧的自链接路径,需要重新定位。

这就是 tail 可以滞后的原因:元素加入的正确性并不依赖 tail 与最后一个节点实时完全一致。把“tail 更新成功”当成入队唯一成功点,会误解源码中本来允许的竞争情形。

offer() 返回 true,也不代表元素仍然等待消费。另一个线程可能在 offer() 返回之前,就把这个已经成功发布的节点消费掉。线程安全保证操作关系正确,不保证调用返回时外部世界停在某个固定状态。

2. poll:原子取得并逻辑删除成员

源码定位:ConcurrentLinkedQueue.java:383–402。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
// 只有成功把非 null 的 item 改为 null 的线程才能取得该成员。
public E poll() {
restartFromHead: for (;;) {
for (Node<E> h = head, p = h, q;; p = q) {
final E item;
if ((item = p.item) != null && p.casItem(item, null)) {
// 这次成功的 CAS 就是该项从队列移除的线性化点。
if (p != h) // 每次跳两个节点
updateHead(h, ((q = p.next) != null) ? q : p);
return item;
}
else if ((q = p.next) == null) {
updateHead(h, p);
return null;
}
else if (p == q)
continue restartFromHead;
}
}
}

poll() 从 head 开始跳过 item 为 null 的节点,遇到有效数据时,用 CAS 尝试把 item 置为 null。如果成功,当前线程获得这个元素;如果失败,说明其他线程已经改变状态,继续查找即可。

这一步同时完成“取得”和“删除”。两个消费者即使看见相同节点,也只有一个能成功把相同的非 null 值改为 null,从而避免同一个队列成员被重复取出。

如果到达 next 为 null 的末端,且没有找到可用成员,poll() 返回 null。由于用户不能存储 null,这个结果能够明确表示本次没有取得数据。

源码中的 updateHead() 还会尽量推进头部入口,但这与某个消费者是否已经取得元素是两个不同动作。已经成功的 item CAS 不需要等待所有失效节点被物理清除。

1
2
3
4
5
6
7
ConcurrentLinkedQueue<Integer> queue = new ConcurrentLinkedQueue<>();
queue.offer(1);
queue.offer(2);
Integer value;
while ((value = queue.poll()) != null) {
System.out.println(value);
}

这个循环只处理当下能够持续取得的成员,第一次 poll() 返回 null 就结束。如果生产者随后才加入新元素,循环不会自动重新开始。需要长期等待的消费者,应考虑阻塞队列或者明确的外部通知机制,避免随意写成不受控制的忙等循环。

3. peek、contains 与 remove

peek() 会跳过逻辑删除节点并查看第一个有效 item,但不会删除它。源码可能顺便推进 head,所以“没有移除用户成员”并不等于“绝不改变内部辅助指针”。

peek() 得到非 null 后再 poll(),是两个独立操作。另一个消费者可以在两者之间先取走它,不能把这段组合代码当成原子的“看后必取”。

contains() 和 remove(Object) 沿链遍历,按 equals() 匹配用户元素,通常需要 O(n) 的扫描。remove() 找到匹配成员后,先通过 item CAS 完成逻辑删除,再尽力清理链接。

1
2
3
4
ConcurrentLinkedQueue<String> queue = new ConcurrentLinkedQueue<>();
queue.addAll(List.of("A", "B", "A"));
System.out.println(queue.remove("A")); // true:删除一个匹配成员
System.out.println(queue); // [B, A]

队列没有去重逻辑,remove(Object) 也不等于移除所有相等成员。若业务需要按唯一任务编号撤销,需要明确“撤销一个”还是“撤销所有”的语义,并考虑消费者已经取得任务的竞争情形。

4. size 为什么不宜作为频繁控制条件

源码定位:ConcurrentLinkedQueue.java:466–478。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
// 没有精确成员计数器,size 需要扫描当时遇到的有效节点。
public int size() {
restartFromHead: for (;;) {
int count = 0;
for (Node<E> p = first(); p != null;) {
if (p.item != null)
if (++count == Integer.MAX_VALUE)
break; // @see Collection.size()
if (p == (p = p.next))
continue restartFromHead;
}
return count;
}
}

size() 不是读取一个 count 字段,而是从有效头节点遍历,统计 item 非 null 的成员。因此它是 O(n) 操作;当规模很大或频繁调用时,成本可能明显超过 offer()/poll()。

并发增删过程中,这次扫描看到的链会持续变化,返回值不一定对应调用者希望的那个统一时刻。即使某次结果恰好准确,随后另一个线程仍然可以立即改变数量。

不能通过 if (queue.size() < limit) queue.offer(task) 实现严格容量控制,因为检查与加入分离,而且多个生产者可以同时通过检查。需要硬容量限制时,应直接选择有界实现,而不是给无界队列外围加一个非原子的 size() 判断。

同样,也不需要先写 if (!queue.isEmpty()) 再 poll()。直接看 poll() 返回的元素,才能避免“检查时非空,消费时已空”的竞争。

addAll() 会先构建一段节点链,再将其连接到队列上。这个具体发布方式不意味着所有 Collection 批量方法都具有事务语义。removeAll()、retainAll()、clear() 或并发遍历与批量更新的组合,不能概括成其他线程只可能看到“全前”或“全后”。

迭代器

ConcurrentLinkedQueue 的迭代器是 弱一致迭代器:允许与队列更新同时进行,不会按照 modCount 机制抛出 ConcurrentModificationException,也不会先复制完整数组作为快照。

创建时会寻找第一个有效节点,把节点保存在 nextNode,把元素引用保存在 nextItem。hasNext() 主要看这个缓存;next() 返回已经缓存的元素,同时寻找下一个有效节点。

源码定位:ConcurrentLinkedQueue.java:775–793。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
// 返回已经缓存的下一元素,同时跳过并尽力清理已删除节点。
public E next() {
final Node<E> pred = nextNode;
if (pred == null) throw new NoSuchElementException();
// assert nextItem != null;
lastRet = pred;
E item = null;

for (Node<E> p = succ(pred), q;; p = q) {
if (p == null || (item = p.item) != null) {
nextNode = p;
E x = nextItem;
nextItem = item;
return x;
}
// 摘除已删除的节点
if ((q = succ(p)) != null)
// 字段仍等于预期值时才完成原子更新。
NEXT.compareAndSet(pred, p, q);
}
}

源码为何要缓存 nextItem?如果 hasNext() 已经确认有成员,另一个消费者随后把对应节点的 item 改为 null,迭代器仍然需要兑现刚刚作出的可读取判断。因此,一个已经被消费者取走的成员仍可能由这个迭代器返回。

这不等于元素被队列重复消费。poll() 通过原子删除取得成员,iterator() 只是读取引用;业务上如果把 for-each 当成消费入口,就会混淆观察与取得所有权这两种操作。

迭代器可能看到某些遍历期间加入的元素,也可能看不到。它保证文档描述的弱一致遍历关系,不能将输出解释为开始时或结束时的完整快照。

1
2
3
4
5
6
ConcurrentLinkedQueue<Integer> queue = new ConcurrentLinkedQueue<>();
queue.addAll(List.of(1, 2, 3));
Iterator<Integer> iterator = queue.iterator();
System.out.println(iterator.next()); // 1
iterator.remove(); // 删除刚返回的那个节点成员
System.out.println(queue); // [2, 3]

本版 Itr.remove() 直接把 lastRet.item 置为 null,然后清空 lastRet。它依赖 volatile 写完成逻辑删除,物理链接可以留给后续遍历处理;不要因为 offer()/poll() 使用 CAS,就断言内部每一次字段更新都必须使用 CAS。

源码定位:ConcurrentLinkedQueue.java:797–803。

1
2
3
4
5
6
7
8
// 迭代器删除的是最近返回的节点,下一次删除前必须再次 next。
public void remove() {
Node<E> l = lastRet;
if (l == null) throw new IllegalStateException();
// 依赖后续遍历重新链接。
l.item = null;
lastRet = null;
}

在还没有 next(),或对同一次 next() 连续 remove() 两次时,会抛出 IllegalStateException。如果其他线程已经删除该节点,迭代器也不会通过 equals() 再找一个相等成员误删。

spliterator() 具有 ORDERED、NONNULL、CONCURRENT 特征。ORDERED 对应队列的遍历顺序,不意味着把一边持续修改的队列交给并行流,就会得到稳定时刻的完整统计。

toArray() 可以方便地把当前遍历所得成员交给后续处理,但复制过程仍然可能遇到并发变化,也不是一个锁住全队列的原子快照。若业务要求“这一批必须与某个事务时刻完全一致”,应该在更上层建立相应的同步或版本边界。

LinkedTransferQueue

基本特性

LinkedTransferQueue 实现 TransferQueue,而 TransferQueue 又扩展 BlockingQueue。它既能像普通无界并发队列一样保存任务,也能让生产者等待任务被接收方取得,是一种兼具缓冲与交接能力的队列。

它不允许 null,允许重复元素,使用链式结构,没有指定容量的构造器。普通 offer()、add()、put() 不等待空余容量;take() 在没有数据时可以等待。新增的 transfer() 与 tryTransfer() 则表达不同强度的交接要求。

操作 当前没有等待接收者时 返回或完成意味着什么
offer() / put() 将元素作为待消费数据节点加入 完成提交,可能尚未被取得
tryTransfer(e) 返回 false,不保留本次元素 成功时已立即匹配接收者
transfer(e) 加入数据节点并等待匹配 等待交接,不等待业务处理完成
tryTransfer(e, timeout, unit) 最多等待指定时间 超时失败的本次交付会被取消
1
2
3
4
5
LinkedTransferQueue<String> queue = new LinkedTransferQueue<>();
System.out.println(queue.tryTransfer("A")); // false,没有接收者
System.out.println(queue.isEmpty()); // true,A 没有留下
queue.offer("B");
System.out.println(queue.poll()); // B,普通 offer 允许缓冲

因此,tryTransfer() 不是一种“提交成功后顺便问问有没有消费者”的 offer()。它一旦失败,本次元素就没有成为待消费任务;调用者必须决定重试、改用 offer(),还是放弃提交。

队列保证元素的 FIFO 关系,但多个竞争线程的实际完成顺序还取决于原子匹配和调度。它没有 ArrayBlockingQueue 那样的公平构造参数,也不能作为限制任务积压数量的有界容器。

结构分析

双重队列(dual queue) 是理解本类的关键:链中不只可以保存生产者的数据节点,也可以保存消费者的请求节点。消费者没拿到数据时,把“我正在等一个元素”登记成请求,后来的生产者可以直接匹配它。

LinkedTransferQueue 的缓冲与直接匹配

本机 JDK 21.0.4 的节点类叫 DualNode,它保存 item、next、waiter、isData 等状态,并实现 ForkJoinPool.ManagedBlocker。不要把其他版本中的旧 Node 定义和 NOW/ASYNC/SYNC/TIMED 常量直接套进本版代码。

isData 区分数据节点与请求节点。未匹配的数据节点有非 null 的 item;被接收后 item 变为 null。未匹配的请求节点 item 为 null;生产者匹配它后,item 被改为交付的对象。

这时,同样的 null 在不同类型节点中表达不同状态,因此判断是否匹配必须结合 isData,不能看到 item 为 null 就一律认为“这是无效节点”。

源码定位:LinkedTransferQueue.java:378–380。

1
2
3
4
// 节点类型与 item 是否有值一起决定匹配状态。
final boolean matched() {
return isData != (item != null);
}

取消请求还有一个专用标记:请求节点可以把 item 设置为节点自身,以区分正常的用户数据。用户无法正常提交这个内部 DualNode 对象,因此自指引用可以充当算法状态标记。

head、tail 不必始终指向第一个和最后一个活跃成员。源码同样允许指针滞后、已经匹配的节点暂留,以及旧头 next 自链接。后续操作会在遍历和清扫时帮助结构前进。

队列同时考虑两件事:匹配必须原子发生,旧节点引用又不能无限保留后续活跃链。前者通过节点 item 的比较交换完成,后者通过推进入口、断开旧路径和取消节点清理完成。

基本操作

1. xfer:四种调用共享同一个核心

源码定位:LinkedTransferQueue.java:574–620。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
// 先匹配互补节点;是否登记和等待,由 ns 所表达的模式决定。
final Object xfer(Object e, long ns) {
boolean haveData = (e != null);
Object m; // 匹配结果;没有匹配时就是 e
DualNode s = null, p; // 已入队节点及其前驱
restart: for (DualNode prevp = null;;) {
DualNode h, t, q;
if ((h = head) == null && // 非立即模式下才初始化
(ns == 0L ||
(h = cmpExHead(null, s = new DualNode(e, haveData))) == null)) {
p = null; // 没有前驱
break; // 否则说明初始化竞争失败
}
p = (t = tail) != null && t.isData == haveData && t != prevp ? t : h;
prevp = p; // 避开已知的自链接尾部路径
do {
m = p.item;
q = p.next;
if (p.isData != haveData && haveData != (m != null) &&
p.cmpExItem(m, e) == m) {
Thread w = p.waiter; // 匹配到互补节点
if (p != h && h == cmpExHead(h, (q == null) ? p : q))
h.next = h; // 推进 head;让旧头自链接
// 通知等待线程继续检查匹配状态。
LockSupport.unpark(w);
return m;
} else if (q == null) {
if (ns == 0L) // 非立即模式下才尝试追加
break restart;
if (s == null)
s = new DualNode(e, haveData);
if ((q = p.cmpExNext(null, s)) == null) {
if (p != t)
cmpExTail(t, s);
break restart;
}
}
} while (p != (p = q)); // 若已自链接则重新开始
}
if (s == null || ns <= 0L)
m = e; // 不等待
else if ((m = s.await(e, ns, this, // 位于或接近头部时先自旋
p == null || p.waiter == null)) == e)
unsplice(p, s); // 已取消
else if (m != null)
s.selfLinkItem();

return m;
}

这个方法的分支较多,但可以顺着一个问题读下来:当前有没有尚未匹配、类型与我相反的节点?有就尝试修改其 item 完成交接;没有,就按照调用模式决定立即失败、加入队列后返回,还是加入后等待。

e 非 null 表示带数据到达,e 为 null 表示请求数据。ns 的含义如下:负值表示异步提交;0 表示只立即尝试;Long.MAX_VALUE 表示无指定期限等待;有限正值表示限时等待。

在匹配分支,p.cmpExItem(m, e) == m 说明比较交换成功,当前线程完成互补匹配。随后 unpark() 等待线程,推进 head,并返回取得的旧状态 m。

没有匹配者时,如果 ns 为 0,当前调用不会建立一个等候将来消费者的数据节点。如果允许登记,就创建 DualNode,原子连接到末端;异步模式到此即可返回,等待模式还需要调用节点 await()。

这里“节点链接成功”与“数据交接成功”不是同一个事件。普通 offer() 只要求前者;transfer() 还关心后者。这一差异解释了为何相同链表结构可以同时提供两种提交语义。

await() 也不是单纯 while 自旋到永远。源码会根据节点位置、线程类型等决定短暂自旋,登记 waiter,检查中断与期限,再通过 park() 或 managedBlock() 等方式等待。

当正常匹配与取消竞争时,谁先成功原子更新节点状态,就决定本次结果。一个失效节点可能暂时仍在链上,但不会因此再次被当作有效任务交给消费者。

2. offer 和 put:允许任务积压

源码定位:LinkedTransferQueue.java:1145–1148。

1
2
3
4
5
// 本类无容量等待,put 采用异步提交模式。
public void put(E e) {
Objects.requireNonNull(e);
xfer(e, -1L);
}

本类的 put() 方法不声明 InterruptedException,也没有先等 notFull 再入队。因为它没有容量限制要等待,方法把数据异步提交给 xfer() 后直接结束。

带超时参数的 offer() 在这里也不使用超时建立容量等待。它与普通 offer() 一样提交元素;不要根据 BlockingQueue 的接口形状,就推断每个实现都会用完这个参数去等空间。

1
2
3
4
5
LinkedTransferQueue<Integer> queue = new LinkedTransferQueue<>();
queue.put(10);
queue.offer(20, 1, TimeUnit.NANOSECONDS);
System.out.println(queue.poll()); // 10
System.out.println(queue.poll()); // 20

提交速度若长期大于消费速度,队列可以不断积压,最后受内存限制。remainingCapacity() 返回 Integer.MAX_VALUE,并不承诺 JVM 真的还有足够内存保存这么多新节点和业务对象。

本类适合需要“普通提交”和“等待交接”两种能力的通道。如果需求首先是严格控制最多积压多少任务,有界 BlockingQueue 更直接。

3. transfer:等待交接

源码定位:LinkedTransferQueue.java:1218–1226。

1
2
3
4
5
6
7
8
9
10
// 交付节点没有立即匹配时,生产者在这个节点上等待。
public void transfer(E e) throws InterruptedException {
Objects.requireNonNull(e);
if (!Thread.interrupted()) {
if (xfer(e, Long.MAX_VALUE) == null)
return;
Thread.interrupted(); // 只可能因中断而失败
}
throw new InterruptedException();
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
LinkedTransferQueue<String> queue = new LinkedTransferQueue<>();
ExecutorService executor = Executors.newSingleThreadExecutor();
try {
Future<?> sent = executor.submit(() -> {
try {
queue.transfer("任务-C");
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new RuntimeException(e);
}
});
System.out.println(queue.take()); // 任务-C
sent.get(); // 交付线程已经可以完成
} finally {
executor.shutdownNow();
}

与 SynchronousQueue 不同,等待 transfer() 的数据节点在本类中属于可观察的待消费数据。它可以计入 size(),也可能被 peek() 和遍历观察到;SynchronousQueue 对外始终提供空集合视图。

还有一个需要结合源码理解的边界:remove(Object)、迭代器删除等操作,也能把待交付的数据节点标成已匹配并唤醒等待的生产者。因此,transfer() 返回不能作为“业务任务已经成功执行”的确认,更不能替代业务层的结果 Future。

如果应用允许管理线程从队列中撤销任务,就应另外记录任务被执行、被取消或执行失败的状态。仅凭交付方法返回,无法恢复完整的业务生命周期。

4. tryTransfer、take 与超时取消

源码定位:LinkedTransferQueue.java:1202–1205。

1
2
3
4
5
// 立即模式只寻找现成接收者,失败时不排入数据节点。
public boolean tryTransfer(E e) {
Objects.requireNonNull(e);
return xfer(e, 0L) == null;
}

限时 tryTransfer() 给消费者一段到达机会,超过时间且没有完成匹配,就取消这次交付。它与普通限时 offer() 的业务意义完全不同:后者在本类中允许缓冲,前者要求交接。

1
2
3
4
LinkedTransferQueue<String> queue = new LinkedTransferQueue<>();
boolean delivered = queue.tryTransfer("限时任务", 10, TimeUnit.MILLISECONDS);
System.out.println(delivered); // false:没有接收者
System.out.println(queue.poll()); // null:失败的任务没有留在队列

take() 与 poll() 也调用 xfer()。空队列上的 take() 使用请求节点等待,普通 poll() 只尝试当前匹配;限时 poll() 给请求节点设置等待期限。null 被保留为请求标志,因而不能作为用户元素加入。

调用者不能通过 hasWaitingConsumer() 的结果保证下一次 tryTransfer() 成功。在查询返回后,该消费者可能已经超时、被中断或与另一生产者匹配。真正的交付结果应以提交操作的返回值为准。

getWaitingConsumerCount() 同样只是等待请求的估计数量。size() 则遍历有效数据节点,不把请求节点算成用户元素,计算过程为 O(n),并发修改时并非稳定快照。

drainTo() 通过反复 poll() 移走当前可取得的数据。被转移的任务可能来自普通 offer(),也可能来自正在 transfer() 的生产者;这些生产者可以随着数据节点被匹配而继续运行。

迭代器

LinkedTransferQueue 的迭代器弱一致,不复制全量成员,不依赖 modCount,也不会将请求节点暴露成用户元素。

Itr 的 advance() 沿着链寻找 isData 为 true 且 item 仍有效的节点,跳过已匹配或取消节点;遇到未匹配的请求段时,不把它当成一批 null 元素返回。自链接的旧路径则会引导它重新从当前 head 查找。

nextItem 缓存保证迭代器能够返回已经确认的下一引用。与其他弱一致遍历一样,缓存建立后该节点可能被别的线程移除,迭代器仍可能返回这个引用。

1
2
3
4
5
6
LinkedTransferQueue<Integer> queue = new LinkedTransferQueue<>();
queue.addAll(List.of(1, 2, 3));
Iterator<Integer> iterator = queue.iterator();
System.out.println(iterator.next()); // 1
iterator.remove(); // 移除迭代器最近返回的数据节点
System.out.println(queue); // [2, 3]

迭代器删除通过 tryMatchData() 改变节点数据状态,并尝试清理链接。这个方法也会唤醒可能在节点上等待的线程,因此“遍历时删除一个成员”在 transfer() 场景里还可能影响生产者是否继续等待。

源码定位:LinkedTransferQueue.java:704–711。

1
2
3
4
5
6
7
8
9
10
// 删除数据节点也通过匹配状态转换完成,并唤醒可能等待的生产者。
final boolean tryMatchData(DualNode p, Object x) {
if (p != null && p.isData &&
x != null && p.cmpExItem(x, null) == x) {
// 通知等待线程继续检查匹配状态。
LockSupport.unpark(p.waiter);
return true;
}
return false;
}

本类的 spliterator() 具有 ORDERED、NONNULL、CONCURRENT 特征,表示有队列顺序、不提供 null 用户成员并支持并发遍历。没有 SIZED 保证,不能把并行遍历开始前的一次 size() 查询,当成遍历过程中准确不变的成员总数。

这里需要再次区分三个结果:offer() 完成提交,transfer() 等待交付状态完成,iterator() 观察有效成员。它们都不能自动说明业务任务是否成功;把这三层语义分清,才能正确选择等待、缓冲和确认机制。

ReferenceQueue

基本特性

part1 提到 ReferenceQueue 时,说到它经常出现在 GC 相关代码里。这里先把边界说清楚:ReferenceQueue 虽然名字里有 Queue,但它不是 java.util.Queue 的实现类。

它位于 java.lang.ref,类声明没有 implements Queue,也没有继承 AbstractQueue。它的用途是接收引用对象的入队通知,而不是保存普通待处理业务元素。

我们通常把 WeakReference、SoftReference 或 PhantomReference 与它关联。引用处理达到相应条件后,引用对象可以进入这个队列,使用者据此清理外部记录或安排后续工作。

最关键的一点是:入队的是 Reference 对象,不是它曾经指向的原对象。 队列返回 Reference<? extends T>,不能把它当作 T 的普通队列来取得原业务对象。

1
2
3
4
5
6
7
ReferenceQueue<Object> queue = new ReferenceQueue<>();
Object value = new Object();
WeakReference<Object> reference = new WeakReference<>(value, queue);
System.out.println(reference.enqueue()); // true,显式入队
System.out.println(queue.poll() == reference); // true,返回引用对象
System.out.println(reference.get()); // null,JDK 21 的 enqueue 会清除指向关系
Reference.reachabilityFence(value); // 这里仍保留原对象的强可达性

这个例子不用 System.gc(),而是验证明确的手动入队行为。原对象仍强可达,引用也可以显式入队,所以观察到队列成员,不能无条件证明原对象已经被 GC 回收。

ReferenceQueue 的公开 API 很小:poll() 立即尝试取出,remove() 等待引用对象入队,remove(timeout) 等待至超时。它没有 offer()/add()、peek()、公开 size(),也没有迭代器。

它的内部入队、出队和等待有同步协调,但 API 与 BlockingQueue 也不同。例如 remove() 的 timeout 使用毫秒,0 表示无限等待,不是立即失败。

结构分析

JDK 21 的这个类维护 volatile head、queueLength、ReentrantLock lock 和 Condition notEmpty。head 通过引用对象的 next 字段连接起已入队成员,没有为每一个入队引用再包一层独立 Node。

ReferenceQueue 的引用对象、被引用对象与入队链

构造 Reference 时,引用对象保存到其 referent 的特殊引用关系,并记录关联的 queue。但是 queue 不会在构造时就把所有注册的 Reference 都强持有在一个成员列表里。

只有真正入队之后,队列链才指向这份引用对象。所以如果业务依赖某个引用的通知,通常还需要在自己的注册表里保存引用对象本身,防止它在还没通知之前也失去可达性。

保存引用对象,不应该反过来强保存它的 referent。例如给弱引用子类加上一个普通字段并把 referent 再放进去,就会破坏原本希望的弱持有关系。

ReferenceQueue 的 NULL 与 ENQUEUED 是内部哨兵对象,用作 Reference.queue 的状态标记。它们都是不接受再次入队的特殊队列,不是用户可以添加的 null 成员。

NULL 表示没有可用的关联队列路径,例如已从队列摘出的引用;ENQUEUED 表示当前已经进入队列。状态标记与 next 链关系共同防止同一引用反复进入。

queueLength 为内部计数,没有对外公开 size() 方法。不能看到这个字段就自己推导出一套 Queue API,更不应依赖反射读取它进行普通业务判断。

另外,这份源码使用 ReentrantLock 和 Condition,不是旧资料中常见的 synchronized(lock) 加 Object.wait()/notifyAll()。源码片段与本机类文件的私有字段签名一致。

基本操作

关联引用与保持注册信息

构造时指定 queue,即建立引用通知关系。例如可以继承 WeakReference 保存一个与 referent 无关的稳定标识:

1
2
3
4
5
6
7
8
9
10
11
class TaggedReference extends WeakReference<Object> {
final String key;
TaggedReference(String key, Object referent, ReferenceQueue<Object> queue) {
super(referent, queue);
this.key = key; // 保存标识,不在普通字段中再次保存 referent
}
}
ReferenceQueue<Object> queue = new ReferenceQueue<>();
TaggedReference reference = new TaggedReference("entry-1", new Object(), queue);
Set<TaggedReference> registrations = new HashSet<>();
registrations.add(reference); // 注册表保存引用对象,供之后消费通知时移除

这个标识可以帮助清理外部索引。获取通知后,使用者通常从注册表中删除这份 Reference,并更新与标识关联的记录,避免通知处理完之后仍保留无用引用对象。

WeakReference 适合弱持有关系;SoftReference 的清除与内存需求等规则有关,不能把它当成容量可预测的普通缓存;PhantomReference 的 get() 始终为 null,通常用于对象生命周期结束相关的后续安排。

这几种引用的具体回收语义不同,但 ReferenceQueue 自己只负责引用对象的链入、取得与等待,不负责决定 referent 应该何时被垃圾收集器处理。

显式 enqueue 与 clear 的区别

手动入队的公开入口属于 Reference,不属于 ReferenceQueue。

源码定位:Reference.java:483–486。

1
2
3
4
5
// JDK 21 先清除引用指向关系,再尝试把这份 Reference 放入其关联队列。
public boolean enqueue() {
clear0(); // 这里有意调用 clear0() 而不是 clear()
return this.queue.enqueue(this);
}

enqueue() 返回是否实际成功加入。第一次成功入队后,重复调用不会产生两份同一引用对象;引用从队列被消费之后,也不能把它当作普通业务节点反复重新提交。

而 clear() 只清除指向关系,不负责将引用对象加入队列。执行 clear() 后立即 poll(),不能假定必然取得对应通知。

1
2
3
4
5
6
7
8
9
ReferenceQueue<Object> queue = new ReferenceQueue<>();
Object value = new Object();
WeakReference<Object> reference = new WeakReference<>(value, queue);
reference.clear();
System.out.println(reference.get()); // null
System.out.println(queue.poll()); // 没有其他入队动作时:null
System.out.println(reference.enqueue()); // true,显式触发入队
System.out.println(reference.enqueue()); // false,已经入队
Reference.reachabilityFence(value); // 排除对象提前不可达带来的 GC 竞争

这个示例没有要求垃圾收集器在某个时刻运行,行为由明确的方法调用决定。System.gc() 只是请求,不适合拿来写“调用之后一定立即在队列中看到引用”的确定性程序。

内部入队:保存 Reference 链与状态

ReferenceQueue 内部 enqueue() 取得 lock 后,进入 enqueue0():

源码定位:ReferenceQueue.java:87–108。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
// 在持锁状态下检查引用是否仍可入队,链入头部,再更新入队状态并通知。
final boolean enqueue0(Reference<? extends T> r) { // 必须在持锁状态下调用
// 检查取得锁之后这个引用是否已经被入队(甚至已被取出)
ReferenceQueue<?> queue = r.queue;
if ((queue == NULL) || (queue == ENQUEUED)) {
return false;
}
assert queue == this;
// 自连接作为链表结束标记,这样若是 FinalReference 就保持非活跃状态。
r.next = (head == null) ? r : head;
head = r;
queueLength++;
// 先加入链表,*之后*再更新 r.queue,避免与并发的入队检查和
// 快速路径 poll() 竞争;volatile 保证顺序。
r.queue = ENQUEUED;
if (r instanceof FinalReference) {
VM.addFinalRefCount(1);
}
signal();
return true;
}

如果 r.queue 已经是 NULL 或 ENQUEUED,就返回 false,避免重复加入或消费后重新加入。

真正链接时,r.next 指向原 head;原来为空,则让 r.next 指向 r 自己,作为链结束标记。之后 head 改为 r,queueLength 增加,最后 r.queue 设为 ENQUEUED。

这里先链入,再发布队列状态,源码用 volatile 访问顺序协调快速 poll() 与入队状态观察。不能随意把状态与链操作前后对调。

当前实现是头插、头取的链结构,不提供 Queue 接口的 FIFO 契约。 用户不应依赖多个引用对象一定按创建、变弱或回收的先后顺序返回。

对于内部 FinalReference,源码还更新 JVM 的相关统计。这是 JDK 内部引用处理分支,不表示普通用户应该使用内部引用类来创建业务容器。

poll:先检查 head,再在锁内摘取

源码定位:ReferenceQueue.java:179–188。

1
2
3
4
5
6
7
8
9
10
11
12
// 空队列可以先由 volatile head 快速判断,真正摘取仍由锁保护。
public Reference<? extends T> poll() {
if (headIsNull())
return null;
// 先获取保护队列结构的锁。
lock.lock();
try {
return poll0();
} finally {
lock.unlock();
}
}

锁外 head 为空时,返回 null。另一个线程随后才入队,当前 poll() 仍然可以正确返回“此次检查时没有成员”,并不需要等待未来元素。

如果 head 看起来非空,获取锁后调用 poll0(),可能发现其他消费者已经先取走,这时同样返回 null。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
// 先更新引用队列状态,再移动 head;摘出的引用保持自连接标记。
final Reference<? extends T> poll0() { // 必须在持锁状态下调用
Reference<? extends T> r = head;
if (r != null) {
r.queue = NULL;
// 先从链表移除,*之前*更新 r.queue,避免与并发的入队检查和
// 快速路径 poll() 竞争;volatile 保证顺序。
@SuppressWarnings("unchecked")
Reference<? extends T> rn = r.next;
// 把自连接的 next 作为链表结束标记处理。
head = (rn == r) ? null : rn;
// 用自连接而不是置为 null,这样若是 FinalReference 就保持非活跃状态。
r.next = r;
queueLength--;
if (r instanceof FinalReference) {
VM.addFinalRefCount(-1);
}
return r;
}
return null;
}

r.queue 先改为 NULL,再移除链头,与入队时的顺序约束相对应。next 指向 r 自己表示这是最后一个节点,此时新 head 为 null;否则 head 转到后继。

摘出的 r.next 继续自连接,不会保留对剩余队列链的无必要引用。queueLength 随后减少,并返回 r。

返回这份引用对象之后,使用者可以读它携带的标识,但不能自动通过 get() 恢复已经被清除的 referent。对于 PhantomReference,get() 本来就始终为 null。

remove 与超时等待

源码定位:ReferenceQueue.java:154–160。

1
2
3
4
5
6
7
8
// 无限等待路径反复尝试 poll0,没有成员才释放锁等待 notEmpty。
final Reference<? extends T> remove0() throws InterruptedException { // 必须在持锁状态下调用
for (;;) {
var r = poll0();
if (r != null) return r;
await();
}
}

公开 remove() 先获取 lock,调用这个内部循环,再在 finally 释放锁。等待 notEmpty 时锁会被释放,使入队者有机会完成链接和通知。

remove(timeout) 检查 timeout 非负,0 委派给无限等待 remove()。正值在内部记录 System.nanoTime(),等待之后重新 poll0(),并扣减已经经过的毫秒预算。

1
2
3
ReferenceQueue<Object> queue = new ReferenceQueue<>();
System.out.println(queue.remove(10)); // 10ms 等待预算内无成员时返回 null
// queue.remove(0); // 等到有引用入队或线程被中断

等待可被中断,抛 InterruptedException。超时没有硬实时保证,实际返回还受线程调度与重新获取锁影响;引用处理线程何时把通知链入也不是由 timeout 参数决定。

而且,remove() 等待的是“引用对象已入队”,不是“referent 的所有资源已经安全释放”。GC 的引用处理阶段、原对象内存回收,以及业务外部资源的生命周期需要分开理解。

清理通知与外部资源

ReferenceQueue 可以帮助发现弱持有对象不再可用、安排清理注册项,但不应把不可预测的 GC 时机当作唯一的文件、连接等资源释放机制。

资源能够显式关闭时,应在业务完成时使用对应的 close() 协议;引用通知可以作为生命周期补充。处理线程也需要自己的停止、中断和异常恢复机制,ReferenceQueue 没有统一 close()。

Reference.reachabilityFence 可用于确保某对象在指定程序点之前保持强可达,避免某些资源使用代码只使用了对象关联状态、却让对象提前满足不可达条件。它不是取消 GC 的永久保活开关,也不会自动建立所有对象字段的线程同步。

一个常见误区是:“队列还能返回 reference,说明它持有 referent。”队列持有的是引用包装对象,弱、软、虚引用的特殊 referent 关系仍按各自规则处理,两层关系不能混同。

迭代器

这个小节需要明确说明:ReferenceQueue 没有公开 iterator(),也不实现 Iterable。 不能对它直接使用增强 for 或调用 stream(),不能为了统一章节结构而虚构一个 Queue 迭代入口。

如果需要处理当前能够取得的通知,就反复 poll():

1
2
3
4
5
6
7
8
ReferenceQueue<Object> queue = new ReferenceQueue<>();
WeakReference<Object> reference = new WeakReference<>(new Object(), queue);
reference.enqueue();
int processed = 0;
for (Reference<?> r; (r = queue.poll()) != null; ) {
processed++; // 这里按 r 携带的标识清理自己的登记记录
}
System.out.println(processed); // 1

这不是不修改队列的遍历,而是消费式读取:每次 poll() 都从队列移除一份通知。其他线程可以同时入队,循环结束只表示最后一次 poll() 没有取得成员,不保证此后永远不会再有通知。

源码还有包内的 forEach() 辅助诊断方法,但它不是用户可调用的 Iterable 契约。不能把反编译看到的内部方法当成公共 API 使用。

在需要持续监听时,可以用 remove() 等待第一份通知,再用 poll() 批量消费当前积压,降低持续空轮询带来的 CPU 开销。等待退出与异常处理则属于消费线程的业务协议。

ReferenceQueue 与普通阻塞队列之间最根本的区别,也在这里体现:前者管理引用处理通知和引用状态,后者管理用户提交的普通成员与容量、等待条件。名称相似,不意味着 API、顺序和生命周期相同。

Deque

Deque 接口概述与 LinkedList 的队列用法

Deque 继承 Queue,名字来自 double-ended queue,也就是双端队列。它的重点不是排序,而是两个端点都能够插入、查看和删除元素。

1
2
3
4
5
6
Deque<String> deque = new ArrayDeque<>();
deque.offerLast("A");
deque.offerLast("B");
System.out.println(deque.pollFirst()); // A,队列:尾入头出
deque.push("C");
System.out.println(deque.pop()); // C,栈:头入头出

普通 Queue 的 offer()/add() 对应 Deque 的尾部添加,poll()/remove() 对应头部移除,peek()/element() 对应头部查看;push() 对应 addFirst(),pop() 对应 removeFirst()。

JDK 21 中,Deque 还继承 SequencedCollection,可以使用 reversed() 查看反向视图。反向视图共享原集合成员,首尾操作会换成对应的另一端;不同实现是否支持修改、怎样保证并发操作,要继续看实际实现。

LinkedList 在 part1 已经完整分析过,这里补足一个容易混淆的队列入口:

源码定位:LinkedList.java:683–686。

1
2
3
4
5
// 队列为空时返回 null,否则摘下首节点。
public E poll() {
final Node<E> f = first;
return (f == null) ? null : unlinkFirst(f);
}

因此,LinkedList.poll() 不是直接调用“空集合抛异常”的公开 removeFirst(),而是先判断 first,再使用 unlinkFirst()。两者最终都能删除首节点,但失败行为不同。

LinkedList 允许 null,首尾指针调整为 O(1),查找指定值或索引则仍可能为 O(n);它没有并发保护,也不会因为实现 Deque 就自动支持阻塞等待。

下面的 ArrayDeque、LinkedBlockingDeque、ConcurrentLinkedDeque,分别解决普通数组双端存储、阻塞式双端协调和并发双端更新问题。

ArrayDeque

基本特性

ArrayDeque 是一个基于可扩容数组的双端队列,实现 Deque,可以在头尾两端加入、查看和删除元素。它允许重复元素,不允许 null,不提供线程安全保证,也没有 put()/take() 这样的条件等待接口。

它可以同时承担队列和栈的角色。作为 FIFO 队列时,从尾部加入、从头部取出;作为 LIFO 栈时,把 push() 与 pop() 都映射到头部。与已经介绍过的 LinkedList 相比,它不为每个元素另外创建双向链表节点。

1
2
3
4
5
6
7
8
9
Deque<String> queue = new ArrayDeque<>();
queue.offerLast("A");
queue.offerLast("B");
System.out.println(queue.pollFirst()); // A:先进先出

Deque<String> stack = new ArrayDeque<>();
stack.push("A");
stack.push("B");
System.out.println(stack.pop()); // B:后进先出

常规头尾操作的摊还时间复杂度为 O(1)。少数加入操作可能触发 O(n) 数组扩容,按一系列操作分析才能说摊还 O(1),不能写成每一次 addLast() 都严格 O(1)。

按值搜索的 contains()、remove(Object)、removeFirstOccurrence()、removeLastOccurrence() 通常为 O(n)。ArrayDeque 没有公开的 get(index) 或 remove(index),不能因为底层是数组,就把它当成带随机索引接口的 List。

构造器中的初始容量只影响最初准备的数组空间,并不是成员数量上限。如果任务越积越多,队列会扩容,最终仍受 JVM 的数组大小和可用内存限制。需要硬容量和阻塞等待时,应选择 ArrayBlockingQueue 等相应实现。

结构分析

主要字段是 Object[] elements、int head、int tail。head 指向逻辑上的第一个成员,tail 指向下一次尾部插入的位置,tail 本身在正常稳定状态下对应空槽,而不是最后一个成员。

默认构造器如下:

源码定位:ArrayDeque.java:180–182。

1
2
3
4
// 默认数组长度为 17,正常状态保留一个空槽,可容纳 16 个成员而不扩容。
public ArrayDeque() {
elements = new Object[16 + 1];
}

这里的 16 是初始可容纳成员数,实际数组长度是 17。不能把 ArrayDeque 的“默认容量 16”直接理解成 new Object[16]。

与较早版本不同,当前实现不要求数组长度必须为 2 的幂。下标回绕用 inc()、dec() 等帮助方法判断边界,不是统一写成 index & (length - 1)。

源码定位:ArrayDeque.java:216–219。

1
2
3
4
5
// 下标到达数组末端后回到 0,不要求 modulus 是 2 的幂。
static final int inc(int i, int modulus) {
if (++i >= modulus) i = 0;
return i;
}
ArrayDeque 的数组回绕与空槽

图中元素的逻辑顺序是 A、B、C、D,尽管对应槽位跨过数组末端。head 到 tail 的循环距离就是成员数量;数组位置的大小顺序不能直接当成队列顺序。

正常状态至少保留一个空槽,所以 head 等于 tail 可以表示空队列,而不与“数组完全被成员填满”长期混淆。头尾加入暂时使二者相等时,方法会立即调用 grow(),恢复足够空间。

这与 ArrayBlockingQueue 不同:后者有 count 字段区分满与空,固定数组的所有槽位都可用于保存成员;ArrayDeque 依靠空槽和 head/tail 的关系表达状态。

不允许 null 也在这里发挥作用。源码将非成员槽位设为 null,用它识别可用位置、清理引用并检测某些异常遍历。允许用户 null 会破坏这套约定,同时让空队列 poll() 的返回值产生歧义。

元素被移除时会清空对应槽位,解除数组对对象的引用。这不保证对象立即被垃圾回收,也不影响调用者或其他容器仍持有的引用。

基本操作

1. addFirst / addLast:从两端插入

源码定位:ArrayDeque.java:283–290。

1
2
3
4
5
6
7
8
9
// 先向前回绕 head,再写入头部;临时占满数组时扩容。
public void addFirst(E e) {
if (e == null)
throw new NullPointerException();
final Object[] es = elements;
es[head = dec(head, es.length)] = e;
if (head == tail)
grow(1);
}

addFirst() 先拒绝 null,将 head 向前移动一格,再把元素写入该槽。如果 head 移动后与 tail 重合,说明本次加入使可用槽位耗尽,需要立即扩容。

尾部操作方向相反:先把元素放在 tail 指向的空槽,再向后移动 tail。

源码定位:ArrayDeque.java:300–307。

1
2
3
4
5
6
7
8
9
// tail 是下一次尾部写入位置,写入后再推进到新的空槽。
public void addLast(E e) {
if (e == null)
throw new NullPointerException();
final Object[] es = elements;
es[tail] = e;
if (head == (tail = inc(tail, es.length)))
grow(1);
}

offerFirst() 和 offerLast() 在本类中调用对应 add() 方法,然后返回 true。因为它没有固定容量限制,所以不会像有界队列一样因“已满”正常返回 false;null 和内存分配失败仍然可以导致异常。

1
2
3
4
5
ArrayDeque<Integer> deque = new ArrayDeque<>(2);
deque.addLast(2);
deque.addFirst(1);
deque.addLast(3); // 初始容量不是上限,可以扩容
System.out.println(deque); // [1, 2, 3]

这里说能够扩容,不代表所有异常发生后都具有事务式回滚保证。特别是内存耗尽等异常,不宜当作日常的容量控制机制。队列业务应在正常资源范围内运行。

2. grow:扩容为何还要移动一段旧元素

源码定位:ArrayDeque.java:139–159。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
// 复制数组后,必要时将跨界的高下标片段平移到新数组末端。
private void grow(int needed) {
// 这段代码特意考虑了溢出
final int oldCapacity = elements.length;
int newCapacity;
// 小数组翻倍,否则增长 50%
int jump = (oldCapacity < 64) ? (oldCapacity + 2) : (oldCapacity >> 1);
if (jump < needed
|| (newCapacity = (oldCapacity + jump)) - MAX_ARRAY_SIZE > 0)
newCapacity = newCapacity(needed, jump);
// 复制数组引用槽位,元素对象本身没有被深拷贝。
final Object[] es = elements = Arrays.copyOf(elements, newCapacity);
// 特殊情况:这里 tail == head 需要进一步区分
if (tail < head || (tail == head && es[head] != null)) {
// 发生回绕,把前一段整体前移到数组末尾
int newSpace = newCapacity - oldCapacity;
System.arraycopy(es, head,
es, head + newSpace,
oldCapacity - head);
for (int i = head, to = (head += newSpace); i < to; i++)
es[i] = null;
}
}

旧数组长度小于 64 时,jump 为旧长度加 2,因此新长度通常为 2 * oldCapacity + 2;较大数组通常增加约一半。还要结合 needed、最大数组长度和溢出检查,不能把所有扩容概括成固定翻倍。

例如默认长度 17 的数组,在普通逐个加入触发扩容时通常增长到 36。增长的是内部数组长度,而正常稳定状态的可用成员数仍然比数组长度少一个。

Arrays.copyOf() 只把旧槽位复制到相同下标。若原队列跨界,例如有效元素先在高下标段、再绕到低下标段,新增数组长度会改变回绕边界,单纯复制就无法保持新的循环路径。

所以源码把原本从 head 到旧数组末端的那段元素,向新数组末端方向平移 newSpace 个位置,并相应增加 head;低下标段保留原位置。这样,从新的 head 走到新末端,再绕回低下标段,仍然得到原来的逻辑顺序。

旧位置还需要置为 null,防止多余引用留在非成员槽位。这里搬动的是对象引用,没有深拷贝元素,业务对象的身份保持不变。

如果 head 与 tail 重合且 es[head] 非 null,说明处于加入后临时满数组的状态,扩容也按跨界场景处理。这解释了 grow() 中为何不能只判断 tail 小于 head。

addAll() 还可能根据整批加入所需空间提前扩容,因此 needed 不一定总是 1。增长公式必须满足至少新增所需槽位,再结合常规增长幅度选择数组长度。

3. pollFirst / pollLast:删除并清空引用

源码定位:ArrayDeque.java:375–384。

1
2
3
4
5
6
7
8
9
10
11
// null 槽表示本次没有头部元素;非空时清槽并推进 head。
public E pollFirst() {
final Object[] es;
final int h;
E e = elementAt(es = elements, h = head);
if (e != null) {
es[h] = null;
head = inc(h, es.length);
}
return e;
}

pollFirst() 读取 head 的槽位。如果是 null,队列为空,返回 null;否则清空原槽,并将 head 向后移动。pollLast() 则先向前计算 tail 的前一个槽位,再移除尾部成员。

removeFirst() 和 removeLast() 在对应 poll() 返回 null 时抛出 NoSuchElementException。getFirst()/getLast() 与 peekFirst()/peekLast() 之间也有类似差别,只是查看不会移除成员。

1
2
3
4
5
6
ArrayDeque<String> deque = new ArrayDeque<>();
deque.addAll(List.of("A", "B", "C"));
System.out.println(deque.peekLast()); // C,保留成员
System.out.println(deque.pollLast()); // C,移除尾部
System.out.println(deque.pollFirst()); // A,移除头部
System.out.println(deque); // [B]

Queue 的 add()/offer() 映射到尾部加入,remove()/poll() 映射到头部移除,element()/peek() 映射到头部查看。push() 是 addFirst(),pop() 是 removeFirst(),因此空栈 pop() 抛异常,而空队列 poll() 返回 null。

4. 删除中间成员:选择搬动更短的一侧

removeFirstOccurrence() 从头向尾查找第一个 equals() 匹配成员,removeLastOccurrence() 从尾向头查找最后一个。找到槽位后,内部 delete() 再决定如何消除数组中的空洞。

源码定位:ArrayDeque.java:605–638。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
// 比较被删位置两侧成员数,移动较短的一侧以维持连续的循环区间。
boolean delete(int i) {
final Object[] es = elements;
final int capacity = es.length;
final int h, t;
// 待删除元素之前的元素个数
final int front = sub(i, h = head, capacity);
// 待删除元素之后的元素个数
final int back = sub(t = tail, i, capacity) - 1;
if (front < back) {
// 把前面的元素整体前移
if (h <= i) {
System.arraycopy(es, h, es, h + 1, front);
} else { // 回绕
System.arraycopy(es, 0, es, 1, i);
es[0] = es[capacity - 1];
System.arraycopy(es, h, es, h + 1, front - (i + 1));
}
es[h] = null;
head = inc(h, capacity);
return false;
} else {
// 把后面的元素整体后移
tail = dec(t, capacity);
if (i <= tail) {
System.arraycopy(es, i + 1, es, i, back);
} else { // 回绕
System.arraycopy(es, i + 1, es, i, capacity - (i + 1));
es[capacity - 1] = es[0];
System.arraycopy(es, 1, es, 0, t - 1);
}
es[tail] = null;
return true;
}
}

front 表示被删成员前面的成员数,back 表示后面的成员数。前面较少,就把前段向后移一格并推进 head;否则把后段向前移一格并回退 tail。

跨界时不能只执行一次从低到高的线性数组复制,因此源码按跨界片段分别搬动,并处理末端与 0 号槽位之间的衔接。这里仍是 O(n) 的最坏成本,但减少了需要移动的引用数量。

1
2
3
4
5
ArrayDeque<String> deque = new ArrayDeque<>(List.of("A", "B", "A", "C"));
deque.removeLastOccurrence("A");
System.out.println(deque); // [A, B, C]
deque.remove("A"); // 等同于 removeFirstOccurrence
System.out.println(deque); // [B, C]

delete() 返回一个方向标志,告诉迭代器删除是否让后段元素左移。迭代器需要据此修正 cursor,否则可能跳过刚刚搬到当前位置的下一个成员。

clear() 清空有效区间并重置 head、tail,但不会自动把数组缩回默认长度。临时装过大量成员的 deque 即使现在很小,仍可能持有较大的引用数组;若确实要释放这部分数组空间,可以在业务允许时创建新的较小容器接替。

5. JDK 21 的 reversed 视图

Deque 在 JDK 21 中支持 reversed()。ArrayDeque 继承接口提供的反向视图,它共享原容器:视图的第一端对应原容器的最后一端,对视图增删会反映到原容器。

1
2
3
4
5
Deque<Integer> original = new ArrayDeque<>(List.of(1, 2, 3));
Deque<Integer> reversed = original.reversed();
System.out.println(reversed); // [3, 2, 1]
reversed.addFirst(4);
System.out.println(original); // [1, 2, 3, 4]

reversed() 既不复制数组,也不让 ArrayDeque 获得线程安全或固定容量。只想得到独立的反序副本时,应该明确创建新容器并复制视图成员,不能把共享视图当作快照。

迭代器

iterator() 从头到尾,descendingIterator() 从尾到头。普通迭代器保存 cursor、remaining 和 lastRet:cursor 是下一槽位,remaining 是创建时剩余数量,lastRet 记录最近返回槽位以支持 remove()。

源码定位:ArrayDeque.java:695–704。

1
2
3
4
5
6
7
8
9
10
// 逐个读取并回绕 cursor;读取到意外的 null 时由帮助方法发现修改。
public E next() {
if (remaining <= 0)
throw new NoSuchElementException();
final Object[] es = elements;
E e = nonNullElementAt(es, cursor);
cursor = inc(lastRet = cursor, es.length);
remaining--;
return e;
}

这一版没有 ArrayList 那样的 modCount/expectedModCount 字段对。next() 通过 nonNullElementAt() 检查当前位置是否意外变成 null;forEachRemaining() 还会检查剩余区间与 tail 的关系,以及处理期间 tail 是否改变。

因此,它仍是文档所描述的尽力快速失败迭代器,但不是“任何外部修改都必然立刻抛异常”。例如一些改动可能不破坏当前检查所依赖的槽位或端点,从而未被这次调用发现。不能用能否抛异常判断代码是否线程安全。

1
2
3
4
5
6
7
8
ArrayDeque<Integer> deque = new ArrayDeque<>(List.of(1, 2, 3, 4));
Iterator<Integer> iterator = deque.iterator();
while (iterator.hasNext()) {
if (iterator.next() % 2 == 0) {
iterator.remove(); // 用当前迭代器支持的删除方式
}
}
System.out.println(deque); // [1, 3]

remove() 调用 delete(lastRet),再按移动方向修正 cursor,并将 lastRet 清为 -1。反向迭代器需要作相反方向的修正,因此它重写了 next() 与 postDelete() 等方法。

还没有 next(),或同一次 next() 后连续 remove() 两次,会抛 IllegalStateException。遍历期间通过 deque.remove() 修改底层容器,不属于这个迭代器自行协调的删除方式。

spliterator() 是晚绑定、尽力快速失败的遍历器,具有 ORDERED、SIZED、SUBSIZED、NONNULL 特征。顺序对应从头到尾,大小可以在没有非法并发修改的条件下用于分割;它没有 CONCURRENT 特征。

clone() 会复制引用数组,使两个 deque 的结构修改独立,但元素对象仍共享。对于可变元素,修改对象字段仍可能被两个容器同时观察到,不能把 clone() 当成对象图的深复制。

LinkedBlockingDeque

基本特性

LinkedBlockingDeque 实现 BlockingDeque,而 BlockingDeque 同时提供双端队列和阻塞队列的能力。与只从尾部生产、头部消费的常规队列相比,它允许两端分别执行加入、移除、限时等待及无限期等待。

它基于双向链表,允许重复元素,拒绝 null,支持并发访问。构造器可以指定正的容量上限;默认容量为 Integer.MAX_VALUE,通常按近似无界使用,但仍有逻辑上限和内存约束。

1
2
3
4
5
6
BlockingDeque<String> deque = new LinkedBlockingDeque<>(3);
deque.offerLast("普通任务-A");
deque.offerLast("普通任务-B");
deque.offerFirst("临时任务");
System.out.println(deque.takeFirst()); // 临时任务
System.out.println(deque.takeLast()); // 普通任务-B

“链表”不意味着没有容量限制。容量满时,offerFirst()/offerLast() 返回 false,addFirst()/addLast() 抛 IllegalStateException,putFirst()/putLast() 则等待空余位置。空 deque 上对应的 poll() 返回 null、remove() 抛异常、take() 等待成员。

通过 addLast()/takeFirst() 可以保持常规 FIFO;通过 addFirst()/takeFirst() 则可以使用 LIFO。双端能力本身不会替你选择业务顺序,两个方向混合操作时,必须先确定应用希望保留什么优先关系。

本类没有传入公平参数的构造器。所谓头部与尾部顺序,是成员的逻辑位置;它与多线程争锁的公平性不是同一个概念。

结构分析

Node 保存 item、prev、next,队列保存 first、last、count 和 capacity。与 LinkedBlockingQueue 的单链表不同,两端都能通过各自方向的链接找到相邻成员。

LinkedBlockingDeque 的两端与同一把锁

更关键的区别在锁:LinkedBlockingQueue 使用 putLock 与 takeLock 两把锁,而 LinkedBlockingDeque 使用 一把 ReentrantLock 保护两端和中间链接,搭配 notEmpty、notFull 两个条件。

头尾操作都可能同时影响空队列到单成员队列、单成员到空队列等状态,也会修改同一个 count。当前实现选择统一锁协调结构,而不是为两个端点各分一把互不相干的锁。

count 是受锁保护的普通 int,因此不用 AtomicInteger。它代表真实有效成员数,容量检查和结构修改在同一个临界区内完成。不能单看字段类型,就认为普通 int 在这里存在数据竞争。

notEmpty 供等数据的线程使用,notFull 供等空间的线程使用。成功加入成员时通知 notEmpty,成功移除成员时通知 notFull;await() 会释放锁,线程醒来后重新争锁并再次检查条件。

头部加入的内部方法如下,调用者必须已经持锁:

源码定位:LinkedBlockingDeque.java:209–223。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
// 容量检查、链接调整和 count 更新都在同一把锁保护下完成。
private boolean linkFirst(Node<E> node) {
// assert lock.isHeldByCurrentThread();
if (count >= capacity)
return false;
Node<E> f = first;
node.next = f;
first = node;
if (last == null)
last = node;
else
f.prev = node;
++count;
notEmpty.signal();
return true;
}

空链表加入第一个成员时,first 与 last 都要指向新节点。非空时,新节点连接到原 first,并修改原 first.prev。尾部的 linkLast() 是对称操作。

这里没有 LinkedBlockingQueue 那样长期保留的哑头节点。空 deque 的 first、last 为 null;删除最后一个成员时,需要同时恢复这两个空端点。

基本操作

1. 两端 offer:容量不足就正常失败

普通 offerFirst()/offerLast() 创建新 Node,获得锁,再调用 linkFirst()/linkLast()。只要 count 未到 capacity,就调整链接并返回 true;否则返回 false,不等待别的线程将来腾出空间。

1
2
3
4
5
LinkedBlockingDeque<Integer> deque = new LinkedBlockingDeque<>(2);
System.out.println(deque.offerLast(1)); // true
System.out.println(deque.offerFirst(2)); // true
System.out.println(deque.offerLast(3)); // false
System.out.println(deque); // [2, 1]

这里“不等待空间”不等于完全不可能等待:普通 offer() 仍要获取结构锁,若其他线程占有锁,当前线程可能在争锁阶段停下来。判断一种方法是否等待,需要分清争锁与等待队列条件这两个阶段。

Queue 的 offer()/add() 映射到尾端,poll()/remove() 映射到头端。BlockingQueue 的 put() 映射到 putLast(),take() 映射到 takeFirst()。只有主动使用双端方法,才表达不同端点策略。

2. putFirst / putLast:队列满时等待

源码定位:LinkedBlockingDeque.java:382–393。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
// 当前版本先用 lock 获取锁,真正满队列时在 notFull 上可中断等待。
public void putLast(E e) throws InterruptedException {
if (e == null) throw new NullPointerException();
Node<E> node = new Node<E>(e);
final ReentrantLock lock = this.lock;
// 先获取保护队列结构的锁。
lock.lock();
try {
while (!linkLast(node))
// 队列满时释放锁等待,返回后重新检查条件。
notFull.await();
} finally {
lock.unlock();
}
}

while 中每次都尝试 linkLast()。如果容量满,await() 释放锁等待;醒来后再次尝试,直到新节点确实加入。不能把 while 改成 if,因为通知不等于空位已经为本线程预留。

这里要严格按当前源码分析:putFirst()/putLast() 使用 lock.lock(),而不是 lock.lockInterruptibly()。中断响应主要体现在进入 Condition.await() 的等待路径,获取锁阶段不按可中断争锁处理。

如果线程已经带着中断标记,但获取锁后立即有空间,当前源码可以直接完成加入,而不经过 await() 抛异常。方法声明 InterruptedException,并不意味着每条分支都必然先检查中断。

这属于当前实现的具体行为,业务取消逻辑仍应按自身需求判断中断状态,不宜依赖这种差异来绕过取消。下面的限时方法则使用不同的获取锁方式。

源码定位:LinkedBlockingDeque.java:422–439。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
// 限时方法在获取锁时也响应中断,并逐次更新等待预算。
public boolean offerLast(E e, long timeout, TimeUnit unit)
throws InterruptedException {
if (e == null) throw new NullPointerException();
Node<E> node = new Node<E>(e);
long nanos = unit.toNanos(timeout);
final ReentrantLock lock = this.lock;
// 获取锁时也响应线程中断。
lock.lockInterruptibly();
try {
while (!linkLast(node)) {
if (nanos <= 0L)
return false;
nanos = notFull.awaitNanos(nanos);
}
return true;
} finally {
lock.unlock();
}
}

awaitNanos() 返回剩余预算,循环中持续更新,避免每次被唤醒都重新等完整 timeout。预算不足时返回 false;如果在临界检查中已有空间,则可以加入并成功返回。

阻塞方法和超时方法都不具备“精确在某个纳秒返回”的调度承诺。唤醒之后仍需争锁和获得 CPU 时间,超时控制的是等待策略。

3. takeFirst / takeLast:空队列时等待

源码定位:LinkedBlockingDeque.java:479–490。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
// 尝试移除头部;没有成员时在 notEmpty 上等待。
public E takeFirst() throws InterruptedException {
final ReentrantLock lock = this.lock;
// 先获取保护队列结构的锁。
lock.lock();
try {
E x;
while ( (x = unlinkFirst()) == null)
// 暂时没有可取元素,释放锁等待生产者通知。
notEmpty.await();
return x;
} finally {
lock.unlock();
}
}

takeFirst() 不先读 size() 再删除,而是在锁内调用 unlinkFirst()。只有真的没有成员才进入 await(),醒来后重复检查,因此多个消费者不会因一次通知重复取得同一个节点。

takeFirst()/takeLast() 同样使用 lock.lock() 获取锁,限时 pollFirst()/pollLast() 才用 lockInterruptibly()。与上面的 put() 一样,应分别看获取锁与等待条件的源码,不能将同一个结论套给所有阻塞方法。

源码定位:LinkedBlockingDeque.java:247–264。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
// 清空被移除节点的 item,修正端点并通知等待容量的线程。
private E unlinkFirst() {
// assert lock.isHeldByCurrentThread();
Node<E> f = first;
if (f == null)
return null;
Node<E> n = f.next;
E item = f.item;
f.item = null;
f.next = f; // 帮助 GC
first = n;
if (n == null)
last = null;
else
n.prev = null;
--count;
notFull.signal();
return item;
}

unlinkFirst() 返回被删 item,把旧 first.item 置为 null,令旧 first.next 指向自身,再更新 first。若没有后继,last 也变为 null;否则新 first.prev 需要置为 null。

旧头自链接一方面减少它对后续链的保留,另一方面供还持有旧节点的迭代器识别路径失效并重新定位。这里不是把业务成员循环链接成环形队列。

1
2
3
4
5
6
LinkedBlockingDeque<String> deque = new LinkedBlockingDeque<>(2);
deque.putLast("A");
deque.putFirst("B");
System.out.println(deque.takeLast()); // A
System.out.println(deque.takeFirst()); // B
System.out.println(deque.pollFirst()); // null

如果希望终止等待,不应仅调用 clear()。clear() 移除成员并通知等空间的生产者,却没有关闭通道的语义,空队列中的消费者仍可能继续等待。应用应通过中断或明确定义的结束消息协调线程生命周期。

4. 中间删除与重复成员

removeFirstOccurrence() 从 first 向后查找,removeLastOccurrence() 从 last 向前查找,使用 equals() 匹配。这是 O(n) 扫描;找到节点后的链接调整本身才是 O(1)。

源码定位:LinkedBlockingDeque.java:291–309。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
// 内部节点删除连接前后邻居,保留旧节点链接给已经到达它的迭代器使用。
void unlink(Node<E> x) {
// assert lock.isHeldByCurrentThread();
// assert x.item != null;
Node<E> p = x.prev;
Node<E> n = x.next;
if (p == null) {
unlinkFirst();
} else if (n == null) {
unlinkLast();
} else {
p.next = n;
n.prev = p;
x.item = null;
// 不要改动 x 的链接,迭代器可能仍在使用它们。
--count;
notFull.signal();
}
}

删除首尾节点可以复用专门的 unlink() 方法,中间节点则让 p.next 指向 n、n.prev 指向 p,再把 x.item 清空、减少 count 并通知 notFull。

源码特别保留中间节点 x 的 prev/next,原因是有迭代器可能已经持有这个节点,需要利用原链接继续前进。物理上与主链脱离,并不代表必须立即把所有字段都清成 null。

1
2
3
4
LinkedBlockingDeque<Integer> deque = new LinkedBlockingDeque<>(5);
deque.addAll(List.of(1, 2, 1, 3));
deque.removeLastOccurrence(1);
System.out.println(deque); // [1, 2, 3]

这个按值删除不与消费线程形成跨操作事务。如果另一个线程已经取得相同任务,队列中删除失败不能说明任务不存在于整个应用;队列成员状态与正在执行的任务状态要分开管理。

5. drainTo、容量查询与反向视图

drainTo() 从 first 开始移走成员,最多移走指定数量,并通知等空间的线程。它不根据调用者过去使用过 takeLast(),就自动改成从尾端批量取得。

1
2
3
4
5
6
LinkedBlockingDeque<Integer> deque = new LinkedBlockingDeque<>(4);
deque.addAll(List.of(1, 2, 3));
List<Integer> batch = new ArrayList<>();
System.out.println(deque.drainTo(batch, 2)); // 2
System.out.println(batch); // [1, 2]
System.out.println(deque); // [3]

目标集合的 add() 失败或抛异常时,不能要求整个操作全部回滚。调用方应选择合适的接收容器,避免在持锁批量转移过程中触发长时间、复杂或重新进入相关队列的回调。

size() 与 remainingCapacity() 在锁下读取 count,可以得到本次读取的数量;但离开锁后其他线程会立刻改变它们,因此不能通过先查询剩余容量、再 put() 来预订空位。offer() 的原子加入结果才是是否得到空间的直接依据。

JDK 21 的 reversed() 由 Deque 提供共享反向视图,视图方法映射到原 deque 的另一端。返回类型是 Deque,并不额外暴露 BlockingDeque 的 putFirst()/takeLast() 等阻塞扩展;需要阻塞双端操作时应直接使用原 BlockingDeque 对象。

迭代器

iterator() 从 first 向后,descendingIterator() 从 last 向前,两者共享 AbstractItr 的遍历逻辑。它们弱一致,不持有一把锁覆盖整个 for-each,也不创建全量快照。

创建和推进到下一节点时,会短暂获取同一把结构锁。nextItem 缓存确保已经找到的下一元素能够返回;推进中跳过 item 为 null 的已删除节点,遇到自链接旧端点时则从相应当前端点重新寻找。

源码定位:LinkedBlockingDeque.java:1086–1104。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
// 返回缓存元素,并在短暂持锁期间寻找下一个有效节点。
public E next() {
Node<E> p;
if ((p = next) == null)
throw new NoSuchElementException();
lastRet = p;
E x = nextItem;
final ReentrantLock lock = LinkedBlockingDeque.this.lock;
// 先获取保护队列结构的锁。
lock.lock();
try {
E e = null;
for (p = nextNode(p); p != null && (e = p.item) == null; )
p = succ(p);
next = p;
nextItem = e;
} finally {
lock.unlock();
}
return x;
}

因此,线程安全的容器不一定提供完全不加锁的迭代器,也不一定提供一个贯穿全遍历的原子快照。本类采用的是逐步持锁、允许中途并发修改的方式。

1
2
3
4
5
6
LinkedBlockingDeque<Integer> deque = new LinkedBlockingDeque<>();
deque.addAll(List.of(1, 2, 3));
Iterator<Integer> iterator = deque.descendingIterator();
System.out.println(iterator.next()); // 3
iterator.remove();
System.out.println(deque); // [1, 2]

iterator.remove() 在锁下检查 lastRet 节点是否仍有 item,有则 unlink();如果已被另一线程删除,不会按 equals() 去删另一个相等成员。调用后清空 lastRet,防止同一次 next() 重复删除。

弱一致表示遍历可以与更新共同进行,不抛出用于检查外部结构修改的 ConcurrentModificationException;它不意味着遍历一定包含所有新成员,也不意味着缓存引用随后仍属于队列。

spliterator() 具有 ORDERED、NONNULL、CONCURRENT 特征。拆分时会在锁下收集部分成员到数组,再交给子遍历器;这只是每次拆分取得一批引用,并不把整个源容器变成固定快照。

如果要严格固定一个批次供业务处理,可以通过明确的转移操作取得成员,并在更上层定义批次边界。把一个弱一致 for-each 的输出当作所有任务的最终清单,会遗漏并发场景里的状态变化。

ConcurrentLinkedDeque

基本特性

ConcurrentLinkedDeque 是一个基于双向链式结构的并发双端队列,实现 Deque。它允许重复元素,拒绝 null,没有指定容量的构造器,并且没有 BlockingDeque 的 put()/take() 和限时等候能力。

它可以在两端并发加入和取出,适合需要线程安全双端操作、又不需要由容器提供等待数据或空间条件的场景。空 deque 上 pollFirst()/pollLast() 返回 null,而不会等待生产者将来提交任务。

1
2
3
4
5
6
7
Deque<String> deque = new ConcurrentLinkedDeque<>();
deque.offerLast("A");
deque.offerFirst("B");
deque.offerLast("C");
System.out.println(deque.pollFirst()); // B
System.out.println(deque.pollLast()); // C
System.out.println(deque); // [A]

Queue 形式的方法仍是尾部加入、头部取出,push()/pop() 则使用头部形成栈。双端能力让应用可以选择方向,但不会自动让任务顺序变成优先级顺序,也不会构成公平的线程调度策略。

与 LinkedBlockingDeque 的差异集中在协调方式与 API:后者有容量、锁和 Condition 等待;本类的核心结构更新通过原子操作推进,无固定容量,也不提供等候方法。不能仅凭类名都有 LinkedDeque,就把两者替换而忽略背压需求。

单次方法线程安全并不让任意组合都原子。例如先 peekLast() 再 pollLast(),另一个线程可以在中间取走原尾成员;先检查 size() 再加入,也不能建立严格的容量限制。

结构分析

每个 Node 保存 volatile 的 prev、item、next,队列保存支持原子更新的 head 与 tail。内部也存在没有用户 item 的辅助节点,所以不是每个物理节点都对应一个有效成员。

ConcurrentLinkedDeque 的双向链接与逻辑删除

与 ConcurrentLinkedQueue 类似,head、tail 可以适度滞后于真正的有效端点。算法从入口沿 prev/next 寻找当前物理端点,再尝试将新节点连接到那个位置。

为什么不要求每次都把两个入口立即调整到最准确的位置?因为这些字段是遍历入口,用户成员的真正加入由端点链接 CAS 决定;把入口更新作为必要条件,会增加竞争,而不会提供相应的业务价值。

删除也分成逻辑删除和物理清理。item 从用户对象原子地变为 null,就说明该成员已被移除;之后 unlink() 才调整前后链接,跳过失效节点并整理端点。

双向结构让清理比单链队列更复杂:一个节点可能从 next 方向不再可达,但仍被另一个节点的 prev 引用。源码需要寻找活跃的前驱和后继,分别修正两个方向,并避免误断正在参与竞争的端点路径。

此外还有 PREV_TERMINATOR、NEXT_TERMINATOR 之类的内部终止标记,以及某些方向的自链接。这些是算法协调与回收辅助状态,不是用户可以加入的 null 元素,也不是把 deque 变成循环链表的公开行为。

创建节点时的 item 写入值得单独看一下:

1
2
3
4
5
6
// 新节点尚未发布,先初始化 item,再借助链接 CAS 的发布关系让其他线程可见。
static <E> Node<E> newNode(E item) {
Node<E> node = new Node<E>();
ITEM.set(node, item);
return node;
}

ITEM.set() 是新节点发布前的初始化,后续原子链接建立可见性关系。不能因为字段声明为 volatile,就断言每个 VarHandle 访问都必须采用最强的访问模式;具体模式应结合发布路径分析。

队列对入队前操作与随后在另一线程中访问或移除元素规定了 happens-before 关系,但只发布引用,不深拷贝任务。任务提交后继续修改其字段,仍需要业务自己设计同步关系。

基本操作

1. linkFirst:找到前端并原子连接

源码定位:ConcurrentLinkedDeque.java:314–341。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
// 成功修改前端节点的 prev 才是新元素加入头部的时刻。
private void linkFirst(E e) {
final Node<E> newNode = newNode(Objects.requireNonNull(e));

restartFromHead:
for (;;)
for (Node<E> h = head, p = h, q;;) {
if ((q = p.prev) != null &&
(q = (p = q).prev) != null)
// 每跳两下检查一次 head 是否已更新。
// 若 p == q,就确定改为从 head 跟进。
p = (h != (h = head)) ? h : q;
else if (p.next == p) // PREV_TERMINATOR
continue restartFromHead;
else {
// p 是首节点
NEXT.set(newNode, p); // 顺带完成 CAS
if (PREV.compareAndSet(p, null, newNode)) {
// 这次成功的 CAS 就是线性化点:
// e 由此成为双端队列中的元素,newNode 也由此“生效”。
if (p != h) // 每次跳两个节点;失败也没关系
HEAD.weakCompareAndSet(this, h, newNode);
return;
}
// CAS 竞争输给其他线程,重新读取 prev
}
}
}

方法首先拒绝 null 并创建新节点。从观察到的 head 开始向 prev 方向走,找到 prev 为 null 的物理前端节点 p,将新节点的 next 预先设为 p。

随后尝试把 p.prev 从 null CAS 为新节点。这个成功的原子操作让新节点成为可达、有效的 deque 成员,也是源码注释明确指出的线性化点。

更新 HEAD 则是辅助后续定位,它使用弱 CAS,失败也不影响已经完成的加入。另一个线程可能已经把 head 推到更合适的位置,无需为保持入口“立即准确”反复争用。

如果寻找路径时遇到终止标记或失效路径,方法重新从当前 head 出发。CAS 失败时重新读状态,而不是覆盖别的线程刚刚设置的前端链接。

对称地,linkLast() 向 next 方向查找尾端,先设置新节点的 prev,再 CAS 原尾端的 next。它们不需要持有一把覆盖整条 deque 的 ReentrantLock。

2. offerFirst / offerLast 与无容量上限

offerFirst()/offerLast() 调用对应 link 方法并返回 true。正常情况下不会因为逻辑容量满而返回 false,但仍可能因 null 或资源不足而异常。

1
2
3
4
5
ConcurrentLinkedDeque<Integer> deque = new ConcurrentLinkedDeque<>();
deque.addLast(2);
deque.addFirst(1);
deque.addLast(3);
System.out.println(deque); // [1, 2, 3]

“无界”是没有预设成员数量边界,不能理解成永远不会耗尽内存。生产者长期超过消费者时,每个待消费节点和其引用的业务对象都会占用资源;本类不自动提供暂停生产者的背压机制。

addAll() 会构建一段双向节点链,再把它连接到原尾部。这个具体路径有自己的原子发布动作,但没有把整个 Collection API 转换为事务:clear()、removeAll()、retainAll() 与其他并发操作的交错仍需按各自语义理解。

3. pollFirst / pollLast:只取到自己成功删除的成员

源码定位:ConcurrentLinkedDeque.java:913–932。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
// item 的 CAS 完成逻辑删除,随后 unlink 尽力整理双向链接。
public E pollFirst() {
restart: for (;;) {
for (Node<E> first = first(), p = first;;) {
final E item;
if ((item = p.item) != null) {
// 重新检查线性一致性
if (first.prev != null) continue restart;
// 字段仍等于预期值时才完成原子更新。
if (ITEM.compareAndSet(p, item, null)) {
unlink(p);
return item;
}
}
if (p == (p = p.next)) continue restart;
if (p == null) {
if (first.prev != null) continue restart;
return null;
}
}
}
}

pollFirst() 从 first 找到的物理前端开始,向后跳过 item 为 null 的节点。看到有效 item 时,尝试把它从原值改为 null;成功则清理节点并返回这个 item。

两个消费者观察到同一个成员时,只有一个能成功完成相同的 CAS,因此不会因为共享节点而重复消费同一个入队成员。

pollLast() 从后端向前做对称工作。双端同时消费同一个只剩一项的 deque 时,仍由 item 的原子状态转换决定谁取得数据,而不是两个端点各自都删除一次。

1
2
3
4
ConcurrentLinkedDeque<String> deque = new ConcurrentLinkedDeque<>();
deque.addLast("唯一成员");
System.out.println(deque.pollLast()); // 唯一成员
System.out.println(deque.pollFirst()); // null

物理清理需要照顾前后两个方向。unlink() 会识别被删节点是否位于前端、后端或中间,跳过相邻已删除节点,调用清理前驱/后继的辅助方法,并在适当时推进入口。

因此,这里的双向链表删除不能直接套成单线程 x.prev.next = x.next; x.next.prev = x.prev。并发时这些邻居可能已被另一线程删除或插入新邻居,必须先校验和使用原子更新。

源码选择允许部分失效节点稍后再清理,让成员删除不必等待所有链接都达到最终形态。逻辑删除是否完成与垃圾回收何时回收节点,仍是两个问题。

peekFirst()/peekLast() 查找并返回有效成员,不删除它。随后通过另一次操作试图取得同一成员时,必须接受其他线程可能已经先一步改变端点的结果。

4. 按值删除与搜索方向

removeFirstOccurrence() 从前向后,removeLastOccurrence() 从后向前,遇到 equals() 匹配的有效 item 时尝试逻辑删除。CAS 失败说明该成员可能已被其他线程取走,需要继续按当前状态工作。

1
2
3
4
5
6
ConcurrentLinkedDeque<Integer> deque = new ConcurrentLinkedDeque<>();
deque.addAll(List.of(1, 2, 1, 3));
deque.removeLastOccurrence(1);
System.out.println(deque); // [1, 2, 3]
deque.remove(1); // remove(Object) 删除从前找到的一个匹配成员
System.out.println(deque); // [2, 3]

按值搜索的复杂度通常为 O(n)。双向结构让应用能够选择从哪一端找第一个匹配项,但不会自动提供索引查找或按值的 O(1) 定位。

如果任务对象的 equals() 依赖可变字段,在入队之后改变这些字段,按值删除是否匹配也可能变化。队列只管理对象引用,不能固定住业务对象的相等关系。

5. size、isEmpty 与共享反向视图

size() 通过遍历有效节点计数,是 O(n),且在并发更新期间可能不精确。它没有一个在每次操作中精确同步更新、可直接读取的成员总数。

不要在消费循环里反复写 size() > 0,既增加扫描成本,又不能防止随后竞争失败。直接检查 pollFirst() 或 pollLast() 的结果通常更符合需求。

isEmpty() 可以根据前端有效成员查找判断,不要求做一次完整 size() 扫描;但这个结果离开当前操作后仍可马上变化,因此也不能用于“先检查,后必定取得”的保证。

1
2
3
4
5
6
ConcurrentLinkedDeque<Integer> deque = new ConcurrentLinkedDeque<>();
deque.addAll(List.of(1, 2, 3));
Deque<Integer> reversed = deque.reversed();
reversed.offerFirst(4);
System.out.println(deque); // [1, 2, 3, 4]
System.out.println(reversed.pollLast()); // 1,映射到原 deque 的前端

JDK 21 的反向视图共享原 deque。视图不会把并发结构变成独立数组快照,其方法通过对应的另一端操作作用于同一成员。通过视图组合两次操作,仍然不构成跨方法事务。

如果想长期等待 deque 变为非空,单凭 reversed() 或 iterator() 不能做到;本类没有条件等待方法。应改用 BlockingDeque 或在应用层建立可靠通知,不能依赖不断无间隔轮询来模拟阻塞。

迭代器

普通 iterator() 从前向后,descendingIterator() 从后向前,都继承 AbstractItr。区别主要由 startNode 与 nextNode 决定,其他缓存、返回与删除逻辑共享。

迭代器是弱一致的,不通过 modCount 对比检查外部结构修改,不复制全量成员,也不在整个遍历期间独占容器。它跳过 item 为 null 的节点,并通过 succ()/pred() 处理失效路径与端点标记。

源码定位:ConcurrentLinkedDeque.java:1390–1408。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
// 寻找下一个有值节点并缓存引用,已删除节点不会作为 null 元素返回。
private void advance() {
lastRet = nextNode;

Node<E> p = (nextNode == null) ? startNode() : nextNode(nextNode);
for (;; p = nextNode(p)) {
if (p == null) {
// 可能位于活跃端或 TERMINATOR 节点,两者都可以
nextNode = null;
nextItem = null;
break;
}
final E item;
if ((item = p.item) != null) {
nextNode = p;
nextItem = item;
break;
}
}
}

nextItem 的作用与 ConcurrentLinkedQueue 相同:一旦已经确认可返回的引用,就算其他消费者随后取走它,迭代器仍可能返回缓存值。因此,遍历是观察,不是原子取得任务。

1
2
3
4
5
6
ConcurrentLinkedDeque<Integer> deque = new ConcurrentLinkedDeque<>();
deque.addAll(List.of(1, 2, 3));
Iterator<Integer> iterator = deque.descendingIterator();
System.out.println(iterator.next()); // 3
iterator.remove();
System.out.println(deque); // [1, 2]

源码定位:ConcurrentLinkedDeque.java:1421–1427。

1
2
3
4
5
6
7
8
// 删除最近返回的节点成员,再帮助清理链接,不按 equals 重新搜索。
public void remove() {
Node<E> l = lastRet;
if (l == null) throw new IllegalStateException();
l.item = null;
unlink(l);
lastRet = null;
}

本版 remove() 直接以 volatile 写把 lastRet.item 清为 null,再调用 unlink()。与主路径 poll() 使用 CAS 的方式不同,但仍是对记录的具体节点删除,不会误把另一个相等对象当成迭代器最近返回的成员。

在尚未 next() 或已对同一次 next() 删除过后调用 remove(),会抛 IllegalStateException;next() 没有下一个缓存元素时抛 NoSuchElementException。弱一致并不免除这些迭代器使用规则。

spliterator() 具有 ORDERED、NONNULL、CONCURRENT 特征。ORDERED 表示前向遇到的成员顺序,不表示存在一个跨整个并行处理固定不变的版本。并行操作的结果需要能容忍源 deque 在处理期间持续更新。

toArray() 通过遍历收集成员,便于后续独立使用这批引用,但采集过程不锁住整个结构,也不保证所有引用属于一个严格统一的瞬间。需要严格批次边界时,应使用业务层同步、转移协议或不同的数据模型。

几种实现放在一起看

上面的类都能参与“加入、观察或取得下一个对象”,但排序、容量和等待语义有很大区别。选用时先确定任务应该怎样被处理,再看对应容器能否直接表达这个约束。

实现 主要结构 固定或可设容量上限 取出规则 等待条件 迭代方式
PriorityQueue 数组二叉堆 没有,初始容量不是上限 比较规则决定堆顶 无 尽力快速失败,非排序遍历
ArrayBlockingQueue 固定环形数组 构造时必需指定 FIFO 满时等空间、空时等成员 弱一致
LinkedBlockingQueue 单链表 可设,默认 Integer.MAX_VALUE FIFO 满时等空间、空时等成员 弱一致
PriorityBlockingQueue 可扩容数组二叉堆 没有 比较规则决定堆顶 空时等成员 数组快照,非排序遍历
DelayQueue PriorityQueue 与锁 没有 按期限选堆顶,过期才可被取出 无可到期成员时等待 数组快照,包括未到期成员
SynchronousQueue 双重节点交接结构 没有缓冲容量 公平 FIFO / 默认 LIFO 匹配 等互补交接方 始终为空
ConcurrentLinkedQueue 单向链式结构与 CAS 没有 FIFO 无 弱一致
LinkedTransferQueue 双重链式结构与 CAS 没有 FIFO 数据匹配 可等成员,也可等交付 弱一致,仅返回数据
ArrayDeque 可扩容循环数组 没有 由两端操作决定 无 尽力快速失败,双向
LinkedBlockingDeque 双向链表与一把锁 可设,默认 Integer.MAX_VALUE 由两端操作决定 满时等空间、空时等成员 弱一致,双向
ConcurrentLinkedDeque 双向链式结构与 CAS 没有 由两端操作决定 无 弱一致,双向
ReferenceQueue 引用通知链接与锁 没有普通容量 API 回收/手动入队通知,无 FIFO 承诺 等引用通知 没有公开迭代器

ReferenceQueue 单列是为了说明其用途,不表示它实现了 Queue。它管理 Reference 通知,不能替换保存业务任务的普通队列。

日常单线程栈或队列通常先考虑 ArrayDeque;要从多个线程安全提交和消费且严格控制积压量,考虑明确容量的 BlockingQueue;要按优先级处理,考虑堆式队列;要等期限到达,考虑 DelayQueue;要直接交接或要求生产者等待交付,再区分 SynchronousQueue 与 LinkedTransferQueue。

“线程安全”只保证方法及容器规定的并发语义,“弱一致”或“快照”描述遍历方式,“有界”描述积压约束,几者不能相互推出。尤其是有锁不必然有界,使用 CAS 不必然没有等待方法,数组也不必然支持 List 的下标接口。

一些面试题

Queue 与 Deque 的区别

Queue 主要定义插入、取得队首和移除队首的操作,普通排队队列一般遵循 先进先出(FIFO)。但 FIFO 不是所有 Queue 实现的共同保证,PriorityQueue 按优先级选择队首;Queue 也没有对外承诺可以独立操作两个端点。

而 Deque 是双端队列,在队列的两端均可以插入或删除元素,也可以模拟栈

Queue 扩展了 Collection 的接口,根据 容量受限、队列为空等正常失败状态的处理方式不同,可以分为两类方法: 一种在操作失败后会抛出异常,另一种则会返回特殊值。

Queue 接口 抛出异常 返回特殊值
插入队尾 add(E e) offer(E e)
删除队首 remove() poll()
查询队首元素 element() peek()

而且 Deque 扩展了 Queue 的接口, 增加了在队首和队尾进行插入和删除的方法,同样根据失败后处理方式的不同分为两类:

Deque 接口 抛出异常 返回特殊值
插入队首 addFirst(E e) offerFirst(E e)
插入队尾 addLast(E e) offerLast(E e)
删除队首 removeFirst() pollFirst()
删除队尾 removeLast() pollLast()
查询队首元素 getFirst() peekFirst()
查询队尾元素 getLast() peekLast()

事实上,Deque 还提供有 push() 和 pop() 等其他方法,可用于模拟栈

ArrayDeque 与 LinkedList 的区别

ArrayDeque 和 LinkedList 都实现了 Deque 接口,两者都能模拟成栈和队列的功能,而且,都非线程安全

  • ArrayDeque 是基于可变长的数组和双指针来实现,而 LinkedList 则通过链表来实现。

  • ArrayDeque 不支持存储 null,但 LinkedList 的具体实现允许 null。List 接口并不要求所有实现接受 null,例如不可修改的 List 工厂会拒绝 null。

    ArrayDeque不是不能支持null,但是poll() / peek() 等方法返回 null 表示“队列为空”。如果允许存 null,就无法区分是“空”还是“存了 null”,容易引发歧义或空指针异常。所以索性就不让存 null 了

  • ArrayDeque 按元素搜索及中间移位删除通常为 O(n)。LinkedList 已经定位节点时,链接调整为 O(1);但按索引定位或按值搜索仍需 O(n),不能把完整的中间增删一概写成 O(1)。ArrayDeque 的数组连续保存元素引用,不额外创建每个元素的链表节点;元素对象本身不一定连续分配。

  • ArrayDeque 插入时可能存在扩容过程, 不过均摊后的插入操作依然为 O(1)。虽然 LinkedList 不需要扩容,但是每次插入数据时均需要申请新的堆空间,会增加节点分配与回收开销;实际性能还要结合访问方式与负载衡量。

  • ArrayDeque 未实现 List 接口,也没用 RandomAccess 接口标记,所以不支持随机访问,而 LinkedList 提供按索引访问接口,但定位节点需要遍历,不具有 O(1) 随机访问能力

对于常规的栈和队列用法,ArrayDeque 通常是更合适的默认选择,具体性能仍应结合实际负载。此外,ArrayDeque 也可以用于实现栈

Java 中的 PriorityQueue 是如何实现的?为什么它不支持 null 元素?线程安全吗?

Java 的 PriorityQueue 是基于可变长数组实现的二叉最小堆(min-heap)。堆调整的 offer()/poll() 成本为 O(log n),peek() 为 O(1)。offer() 偶尔还会触发 O(n) 的数组扩容,按一系列操作分析,其摊还成本为 O(log n)。

关于 null 元素,PriorityQueue 不允许插入 null。原因在于它的 poll() 和 peek() 方法在队列为空时会返回 null。如果允许存入 null,就无法区分队列为空还是堆顶元素就是 null,会导致语义歧义,甚至引发逻辑错误或空指针异常。这是一种防御性编程设计。

此外,PriorityQueue 不是线程安全的。如果多线程并发访问,必须手动加锁,或者改用并发版本 PriorityBlockingQueue。

而且,PriorityQueue 既然基于可变长数组实现的,那么它也有初始的内存分配,并且会自动扩容。机制和 ArrayList 很像

PriorityQueue虽然内部是堆结构,但迭代器不保证有序,因为堆只保证堆顶最小,其余元素无全局顺序。

如何用 PriorityQueue 找出一个大数组中最大的 K 个元素?要求时间复杂度尽可能低。

创建一个容量为 K 的 最小堆,也就是 PriorityQueue,然后遍历数组:

  • 如果堆 size() < K,直接加入;
  • 否则,若当前元素 > 堆顶(即当前 K 个中的最小值),则弹出堆顶,加入当前元素。

遍历结束后,堆中就是最大的 K 个元素

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
import java.util.*;

public class TopK {
public static List<Integer> findTopK(int[] nums, int k) {
if (k <= 0 || nums == null || nums.length == 0)
return new ArrayList<>();

// 最小堆:堆顶是当前 K 个中的最小值
PriorityQueue<Integer> minHeap = new PriorityQueue<>(k);

for (int num : nums) {
if (minHeap.size() < k) {
minHeap.offer(num);
} else if (num > minHeap.peek()) {
minHeap.poll(); // 弹出最小的
minHeap.offer(num); // 加入更大的
}
}

// 返回结果(顺序不重要,或可排序)
return new ArrayList<>(minHeap);
}

public static void main(String[] args) {
int[] arr = {3, 2, 1, 5, 6, 4};
System.out.println(findTopK(arr, 2)); // 这里是 [5, 6];更多元素时内部数组不保证整体排序
}
}

什么是 BlockingQueue?其实现类有哪些?

BlockingQueue (阻塞队列)是一个接口,继承自 Queue。BlockingQueue阻塞的原因是其支持当队列没有元素时一直阻塞,直到有元素;还支持如果队列已满,一直等到队列可以放入新元素时再放入。

BlockingQueue 常用于生产者-消费者模型中,生产者线程会向队列中添加数据,而消费者线程会从队列中取出数据进行处理

BlockingQueue

Java 中常用的阻塞队列实现类有以下几种:

  1. ArrayBlockingQueue:使用数组实现的有界阻塞队列。在创建时需要指定容量大小,并支持公平和非公平两种方式的锁访问机制。
  2. LinkedBlockingQueue:使用单向链表实现的可选有界阻塞队列。在创建时可以指定容量大小,如果不指定则默认为Integer.MAX_VALUE。和ArrayBlockingQueue不同的是, 它仅支持非公平的锁访问机制。
  3. PriorityBlockingQueue:支持优先级排序的无界阻塞队列。元素必须实现Comparable接口或者在构造函数中传入Comparator对象,并且不能插入 null 元素。
  4. SynchronousQueue:同步队列,是一种不存储元素的阻塞队列。阻塞的 put()/take() 需要与另一端的取得/交付操作匹配;普通 offer()/poll() 在没有立即匹配者时直接返回失败或 null,而不是一律等待。因此,SynchronousQueue通常用于线程之间的直接传递数据。
  5. DelayQueue:延迟队列,poll()/take() 只有在堆顶到期后才能取得元素;remove(Object)、clear() 等管理性删除不受这个到期条件限制。

ArrayBlockingQueue 和 LinkedBlockingQueue 有什么区别?

ArrayBlockingQueue 和 LinkedBlockingQueue 是 Java 并发包中常用的两种阻塞队列实现,它们都是线程安全的。不过,不过它们之间也存在下面这些区别:

  • 底层实现:ArrayBlockingQueue 基于数组实现,而 LinkedBlockingQueue 基于链表实现。
  • 是否有界:ArrayBlockingQueue 是有界队列,必须在创建时指定容量大小。LinkedBlockingQueue 创建时可以不指定容量大小,默认上限为 Integer.MAX_VALUE,通常称为近似无界,但仍有逻辑上限,也受到实际内存资源约束。但也可以指定队列大小,从而成为有界的。
  • 锁是否分离: ArrayBlockingQueue中的锁是没有分离的,即生产和消费用的是同一个锁;LinkedBlockingQueue中的锁是分离的,即生产用的是putLock,消费是takeLock,这样可以减少正常入队与出队之间的锁竞争;跨边界通知及某些整体操作仍需要协调两把锁。
  • 内存占用:ArrayBlockingQueue 需要提前分配数组内存,而 LinkedBlockingQueue 则是动态分配链表节点内存。这意味着,ArrayBlockingQueue 在创建时就会占用一定的内存空间,且往往申请的内存比实际所用的内存更大,而LinkedBlockingQueue 则是根据元素的增加而逐渐占用内存空间

PriorityQueue 的初始容量就是元素数量上限吗?堆序能保证遍历顺序吗?

初始容量只是数组最初准备的空间,元素增多后可以扩容。最小堆保证父节点不大于子节点,进而保证堆顶最小,但不保证数组从左到右整体有序。

重复 poll() 可以按优先级取得成员;直接 for-each、toArray() 或 newArrayList<>(queue) 不能被当成有序结果。要在保留原队列的情况下得到排序结果,可以复制后排序,或复制成另一优先队列再逐次 poll()。

优先队列中比较结果为 0 的对象会被去重吗?相同优先级如何保持先后?

不会。PriorityQueue 与 PriorityBlockingQueue 允许重复成员,比较规则主要决定堆位置,不承担 Set 的去重语义。contains()/remove(Object) 依据 equals() 查找,而不是依据比较结果为 0。

同优先级的取出顺序也不自动稳定。如果需要确定顺序,可以为元素增加序号,让比较器先比较优先级,再比较序号,并明确序号的分配时刻与溢出处理。对原始成员入队后随意改写参与比较的字段,不会自动触发重新堆化。

ArrayBlockingQueue 的公平模式保证任务执行顺序吗?

它主要控制竞争内部锁的公平策略,成员仍按照 FIFO 出队。但消费者取到任务后的执行速度、阻塞和线程调度都可能不同,因此先出队的任务不一定先完成。

SynchronousQueue 的公平参数针对等待交接者的匹配次序,同样不能推导业务完成次序。成员顺序、获取机会和业务结果顺序应分别讨论。

size 与 remainingCapacity 能用来预先保证下一次操作成功吗?

不能。查询结果离开当前同步时刻后就可能变化,查询与下一次加入、移除不是一个原子组合。ConcurrentLinkedQueue、ConcurrentLinkedDeque、LinkedTransferQueue 的 size() 还要沿链遍历,通常为 O(n),并发时可能不精确。

应直接使用 offer()/poll() 的结果表达本次成功或失败,用有界容器本身控制容量,而不是给无界队列套上 size() < limit 的外部检查。

SynchronousQueue 的 size 总为 0,put 为什么仍会等待?

size() 是对外的集合视图,不统计内部等待交接节点。put() 要等待互补接收者完成匹配,并不是等待一个容量为 1 的槽位。

普通 offer()/poll() 没有立即匹配者就失败;iterator() 始终为空。当前 JDK 21.0.4 使用 Transferer 复用 LinkedTransferQueue 的节点与等待机制,公平模式采用 FIFO,默认模式使用 LIFO 匹配。

offer、tryTransfer、transfer 的差别是什么?transfer 返回就是任务执行成功吗?

LinkedTransferQueue 的 offer() 允许缓冲任务;立即 tryTransfer() 只匹配已经等待的接收者,失败不保留元素;transfer() 可以登记数据节点并等待交接状态完成。

transfer() 不是业务处理结果确认。消费者拿到对象后尚需执行业务;此外,管理性 remove() 或迭代器删除也可能解除该节点的交付等待。若要确认执行结果,应另外使用 Future、状态记录或应用层协议

DelayQueue 非空,poll 为什么还会返回 null?peek 能看到未到期对象吗?

size() 统计所有登记成员,peek() 可以查看尚未到期的堆顶;poll()/take() 还要检查堆顶 getDelay() 是否不大于 0,因此非空与当前可取出不是同一条件。

remove(Object)、clear()、迭代器删除不要求到期,而 drainTo() 只转移已到期的堆顶成员。不能将所有“移除”方法都概括成必须等待到期。

DelayQueue 为什么只让一个 leader 定时等待?新增更早任务怎么办?

让所有消费者都按相同最早期限定时等待,会产生多余的竞争唤醒。当前实现用 leader-follower 机制,让一个线程负责相应期限,其他线程等待通知。

若新元素成为更早的堆顶,offer() 清空 leader 记录并 signal(),使线程重新确定等待期限。队列不创建业务执行线程,应用仍须自行 take() 并处理任务。

并发队列的迭代器都是快照吗?为什么可能返回已移除的对象?

不是。PriorityBlockingQueue、DelayQueue 使用数组快照;多种链式并发队列以及 ArrayBlockingQueue 使用弱一致遍历;SynchronousQueue 对外迭代器为空。

弱一致迭代器可能缓存已经确认的 nextItem,因此另一个线程删除节点后,迭代器仍可返回缓存引用。快照则固定创建时收集的引用,但可变对象本身没有深拷贝,其字段仍可能改变。

ReferenceQueue 是 Queue 的实现吗?clear、enqueue 与垃圾回收通知有什么区别?

不是。ReferenceQueue 不实现 Queue、Collection 或 Iterable,它用于接收 Reference 对象的通知,没有公开迭代器。

Reference.clear() 清空被引用对象,不负责手动入队;Reference.enqueue() 会清空引用并尝试把引用对象放入其注册队列,因此出队也可能来自手动操作,不能单凭它证明垃圾回收已经完成。通知时机受引用类型和垃圾回收过程影响,示例不应依赖一次 System.gc() 就立即取得结果。