☰
ArrayBlockingQueue源码解析:锁、Condition与生产实践
2026/10/2 3:51:56 网站建设 项目流程

先说一个我自己的真实感受:在 Java 面试里问“ArrayBlockingQueue 用过没”,十个人有九个能背出来“有界阻塞队列、基于数组、FIFO”,但真到线上排查生产问题的时候,能说清楚它内部那套锁和条件是怎么协作的、为什么吞吐量上不去、什么时候该换 LinkedBlockingQueue 的人,我遇到的确实不多。

ArrayBlockingQueue 看起来简单,但它的同步模型其实浓缩了并发编程中相当经典的一整套思路:循环数组、可重入锁、条件队列、等待通知机制。搞懂它,对理解整个 java.util.concurrent 包里的阻塞队列体系都有帮助,而且它也是线程池、消息中间件、本地任务缓冲等场景里最常见的底层组件之一。

这篇不打算写成流水账式的 API 手册,我尽量以“源码 + 运行机制 + 生产踩坑”的视角,把它拆开揉碎。适合正在准备并发相关的同学,也适合那些已经用过它、但想弄明白“为什么有时候加锁了还是会有怪问题”的开发者。

1. 一个锁、一张循环数组:ArrayBlockingQueue 的核心存储模型

在聊 API 之前,建议先把它的“底子”看明白。ArrayBlockingQueue 之所以叫 Array,是因为它内部就是用一个固定长度的 Object 数组来装元素的。这一点决定了它和 LinkedBlockingQueue 最根本的差异:Array 是连续内存、固定容量,Linked 是节点分散、链表式伸缩。

1.1 循环数组:指针怎么往前走

先看三个关键字段的概念模型:

  • items:Object[],存放元素的底层数组
  • takeIndex:下一次 take/peek 时读取的元素下标
  • putIndex:下一次 put/offer 时写入的元素下标
  • count:当前队列中的元素数量

为什么说是“循环”数组?因为 putIndex 和 takeIndex 在走到数组末尾时不会停住,而是通过取模回到 0:

putIndex = (putIndex + 1) % items.length; takeIndex = (takeIndex + 1) % items.length;

举个例子,容量为 4 的队列,放进 a、b、c、d 四个元素后 putIndex 回到 0,这时再 remove 掉队头的 a,takeIndex 从 0 走到 1。此时如果你再 put 一个 e,它会写到下标 0 的位置,把已经被消费掉的“坑位”重新利用起来。

这个设计的价值在于:出队操作不需要像 ArrayList 那样把后面所有元素往前搬移,时间复杂度是 O(1)。只要 count 没满,put 和 take 永远只操作各自的指针,互相之间的“物理位置”没有冲突。

1.2 锁和条件变量:两个条件对应两个方向

存储模型之外,ArrayBlockingQueue 的并发控制模型才是精华。它内部只有一个 ReentrantLock,以及由这把锁派生出的两个 Condition:

final ReentrantLock lock; private final Condition notEmpty; private final Condition notFull; 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(); }

这里有个很值得品味的细节:为什么用两个 Condition,而不是一个?

想象一下只用一把锁 + 一个条件变量的场景:队列满了,生产者在等;队列空了,消费者在等。如果只有一个条件变量,生产者被唤醒时也许是因为消费者取走了元素,但也有可能唤醒的是一个消费者(比如队列空了,消费者也在同一个条件上等待)。这种“信号错乱”需要反复用 while 循环判断,效率不高,而且容易产生无意义的竞争。

两个条件的好处是各等各的:生产者只等“notFull”,消费者只等“notEmpty”。当消费者 take 成功时,准确唤醒一个等待在 notFull 上的生产者;当生产者 put 成功时,准确唤醒一个等待在 notEmpty 上的消费者。这样既避免了惊群,也让语义更清晰。

提示:这两个 Condition 都是基于同一个 ReentrantLock 的,所以“同时只能有一个线程持有锁”的约束并没有变。这个点后面讨论性能时会再展开。

2. 构造参数里的细节:容量、公平锁与初识 onEmpty

构造一个 ArrayBlockingQueue 比很多人想象中更有讲究。构造函数一共有三个重载,最常被忽略的是 fair 这个布尔参数。

2.1 容量为什么必须显式指定

和 LinkedBlockingQueue 不同,ArrayBlockingQueue“天生”必须知道自己的上限。你没法创建一个“不限制容量”的 ArrayBlockingQueue,因为底层数组的大小在构造时就固定了。

容量传 0 会在构造阶段直接抛 IllegalArgumentException,这一点常被用来做参数校验,但实际业务中没人会传 0。比较常见的问题是容量设得过大或过小:设小了,生产者频繁阻塞;设大了,内存占用和 GC 压力上去了,毕竟数组本身是连续分配的一段引用空间。

2.2 fair 参数的代价被低估了

public ArrayBlockingQueue(int capacity, boolean fair) public ArrayBlockingQueue(int capacity, boolean fair, Collection<? extends E> c)

fair 为 true 时,内部 ReentrantLock 走公平模式,等待时间最长的线程优先获得锁。理论上这能避免线程饥饿,但代价非常现实:公平锁需要在线程被唤醒后重新检查队列中是否有更早排队的线程,上下文切换和 CAS 操作的频率会明显上升,吞吐量在竞争激烈时会明显下降。

我自己在做本地任务缓冲中间件压测时对比过:8 个生产者线程 + 8 个消费者线程,队列容量 1024,每轮放 10 万条任务,fair=true 的吞吐大约只有 fair=false 的 60% 左右。除非你有强需求保证某个生产者或消费者不被长时间饿着,否则默认 false 更符合绝大多数生产场景。

2.3 一个容易忽略的初始化逻辑:构造时预放元素

第三个构造函数允许传入一个 Collection,在构造时一次性放入初始元素。它内部会遍历集合并依次调用 enqueue,如果集合元素个数超过容量,会在插入到第 capacity+1 个元素时抛出 IllegalStateException。

这类构造函数实际用得不多,但它有个隐藏特性:传入的集合如果是无序的,队列的初始顺序就完全取决于该集合迭代器的顺序。做测试的时候别想当然认为“我传的 Set 是什么顺序,队列就是什么顺序”。

2.4 性能和内存上的补充说明

对象头 + 数组本身带来的开销:容量为 N 的队列,数组引用占 8N 字节(压缩指针下 4N),加上锁、两个条件、索引等固定开销。别小看这几个字段,当你需要创建几千个队列实例(比如按业务维度隔离队列)时,内存差就不是小数了。

3. 核心 API 的阻塞语义:put/take 与 offer/poll 的取舍逻辑

ArrayBlockingQueue 常用的方法可以按“是否阻塞”分成几组。这一节我们重点把 put/take 和 offer/poll 的机制讲透,顺便解释为什么很多开源框架更喜欢用带超时的 offer/poll。

3.1 put 的内部流程:等待与入队的完整链路

看 put 的源码实现(基于 JDK 8 到 21 的实现思路基本一致):

public void put(E e) throws InterruptedException { checkNotNull(e); final ReentrantLock lock = this.lock; lock.lockInterruptibly(); try { while (count == items.length) notFull.await(); enqueue(e); } finally { lock.unlock(); } }

流程拆开看:

  1. 先检查元素是否为 null,不允许 null 入队
  2. 获取锁,但是用的是 lockInterruptibly,也就是说等待锁的过程中响应中断
  3. 进入临界区后,用一个 while 循环判断队列是否已满。如果满了,就调用 notFull.await() 把当前线程挂起,释放锁
  4. 被唤醒后,重新抢到锁,再次检查 count == items.length。记住:这里必须是 while 而不是 if,原因有两个:一是 wait 可能被虚假唤醒(spurious wakeup),二是可能有多个生产者在等待,某个消费者取走一个元素后,notFull.signal() 只唤醒其中一个,但被唤醒者需要重新确认队列到底还满不满
  5. 队列有空位了,执行 enqueue:
private void enqueue(E x) { final Object[] items = this.items; items[putIndex] = x; putIndex = inc(putIndex); count++; notEmpty.signal(); }

enqueue 本身在锁保护下完成,写入数组、移动 putIndex、count 自增,最后唤醒一个等待中的消费者。

这里有个很多人没想明白的问题:为什么 enqueue 里要调 notEmpty.signal(),而不是 take 方法里自己 signal 自己?其实本质是——入队动作本身意味着“队列非空”这一状态发生了变化,这个状态变化需要通知可能正在等待的消费者。在并发场景下,通知谁、谁去唤醒谁,讲究的是“状态变化由变更方广播”,由 put 通知 notEmpty,由 take 通知 notFull。

3.2 take 的内部流程:与 put 完全对称的镜像

public E take() throws InterruptedException { final ReentrantLock lock = this.lock; lock.lockInterruptibly(); try { while (count == 0) notEmpty.await(); return dequeue(); } finally { lock.unlock(); } }

dequeue 里做的事情是:

private E dequeue() { final Object[] items = this.items; @SuppressWarnings("unchecked") E x = (E) items[takeIndex]; items[takeIndex] = null; takeIndex = inc(takeIndex); count--; notFull.signal(); return x; }

注意这里把 takeIndex 位置的元素置为 null,这一步很关键:如果不置 null,数组里会一直残留已出队元素的引用,导致对象无法被 GC 回收。这在长生命周期的高容量队列上会造成内存泄漏。源码里几乎每个出队路径(remove 方法、removeAt 方法、iterator.remove 等)都会把对应位置清空,这是 JDK 团队对内存安全非常执着的体现。

3.3 offer/poll:阻塞和非阻塞的边界控制

offer 和 poll 本身不阻塞(或只阻塞有限时间),它们在实现上是 put/take 的“非阻塞版本”:

public boolean offer(E e) { checkNotNull(e); final ReentrantLock lock = this.lock; lock.lock(); try { if (count == items.length) return false; else { enqueue(e); return true; } } finally { lock.unlock(); } } public E poll() { final ReentrantLock lock = this.lock; lock.lock(); try { return (count == 0) ? null : dequeue(); } finally { lock.unlock(); } }

这两个方法的特点是不抛 InterruptedException,也不会让线程挂起。它们的价值在于“可控制失败”。举个很常见的生产场景:一个任务提交服务,你希望当队列满了之后能快速返回错误给上游,而不是让上游线程一直阻塞,避免线程堆积。这种场景用 put 就是灾难,用 offer(e, timeout, TimeUnit) 或者 offer(e) 就恰到好处。

带超时的 offer 实现多了一层循环等待:

public boolean offer(E e, long timeout, TimeUnit unit) throws InterruptedException { checkNotNull(e); long nanos = unit.toNanos(timeout); final ReentrantLock lock = this.lock; lock.lockInterruptibly(); try { while (count == items.length) { if (nanos <= 0) return false; nanos = notFull.awaitNanos(nanos); } enqueue(e); return true; } finally { lock.unlock(); } }

值得提醒的是 awaitNanos 返回的是“剩余等待时间”,可能由于各种原因被提前唤醒,所以底层用 while 循环重新判断,并在 nanos <= 0 时返回失败。这种写法是 Java 并发库的模板范式,写自己的等待逻辑时也应该照这个来。

3.4 peek/element 与 remove 的边界行为

  • peek():返回队头元素但不移除,队列为空时返回 null
  • element():peek 的加强版,空队列时抛 NoSuchElementException
  • remove(Object o):遍历查找并移除一个元素,成功返回 true

remove(Object) 的实现比想象中复杂,因为移除的并不总是队头元素,可能出现在数组中间位置。看 removeAt 的实现时,它需要处理“索引回绕”下的元素搬移,最坏情况下要把从被删位置到 takeIndex 之间的所有元素整体前移一位。这绝对是 O(n) 操作,性能不便宜,线上慎用。

4. 源码视角的竞争真相:为什么它只有一个锁,以及 “take 完唤醒 put” 的完整闭环

很多人拿 ArrayBlockingQueue 和 LinkedBlockingQueue 对比后会冒出疑问:ArrayBlockingQueue 只有一个锁,生产者消费者会互相竞争同一把锁,是不是性能上天然吃亏?LinkedBlockingQueue 有 takeLock 和 putLock 两把锁,是不是一定更快?

这个问题的答案比直觉复杂。这一节我们把竞争模型拆开,说说 ArrayBlockingQueue 在什么情况下会暴露性能短板,什么情况下其实足够好。

4.1 两个条件、一个锁:唤醒链的典型路径

完整走一遍 3 个线程的场景:

  • 线程 A:生产者,队列满时阻塞在 notFull
  • 线程 B:消费者,执行 take,成功拿走元素,count 从 max 变成 max-1
  • 线程 C:生产者,马上尝试 put,发现 count < max,直接入队

过程里 B 拿锁 -> 出队 -> notFull.signal() -> 释放锁,A 此时并不会立刻被唤醒继续执行,而是要等 B 释放锁后去竞争同一把锁。也就是说,虽然“notFull 被 signal”了,A 也要参与锁竞争,并且可能输给 C。这没有错,这是正常竞争,代价是 A 可能多等一会。

但问题在另一个方向:ArrayBlockingQueue 生产者消费者共享同一把锁,意味着 put 和 take 之间的并发度不是真正的“读写并行”。如果用 LinkedBlockingQueue,put 锁和 take 锁独立,入队出队可以并行推进,高并发下的整体吞吐确实更容易做高。

4.2 单锁模型为何还能被广泛使用

原因在于真正决定性能的往往不是锁的个数,而是临界区的大小和竞争频率。ArrayBlockingQueue 的临界区极短:入队就是写一个数组槽位、移动索引、count++、signal;出队就是读一个槽位、置 null、移动索引、count--、signal。锁竞争时间被压缩到极短,所以大多数业务场景下单锁并不会成为明显瓶颈。

真正让它吃亏的是“队列本身很短 + 操作频率极高”的组合。比如队列容量只有 1,生产者消费者你来我往,每次操作后 count 快速在 0 和 1 之间跳变,这时所有 put/take 都在同一把锁上排队,性能会显著劣化。反过来,容量充足、操作不极端频繁,ArrayBlockingQueue 的性能表现通常很不错,而且因为它不需要额外创建节点对象,GC 压力和内存碎片比 LinkedBlockingQueue 好。

提示:如果压测发现 ArrayBlockingQueue 成为瓶颈,先检查队列容量设置是否过小,再考虑换结构,别一上来就“无脑换 LinkedBlockingQueue”。容量过小造成的竞争加剧,换任何队列都救不了。

4.3 弱一致迭代器:读快照不是实时快照

ArrayBlockingQueue 的迭代器是弱一致性的(weakly consistent),意思是迭代器创建后,不会抛 ConcurrentModificationException,但也不保证能实时看到后续新增的元素。

Iterator<E> it = queue.iterator(); queue.offer(x); while (it.hasNext()) { // 不一定能看到 x }

它底层维护了一个 nextItem 引用和一套独立的索引推算逻辑,利用内部数组和 count 在创建时刻的“快照”来推导遍历位置。遍历期间如果发生了元素搬移(比如 removeAt 移动了元素),迭代器可能漏掉或重复遍历部分元素。实际开发中我基本只在打印队列状态、做简单统计时用它的迭代器,绝不依赖它做精确遍历和删除。

4.4 批量操作的便捷与风险:drainTo 的正确姿势

drainTo 是处理批量消费最顺手的方法:

public int drainTo(Collection<? super E> c, int maxElements) { checkNotNull(c); if (c == this) throw new IllegalArgumentException(); if (maxElements <= 0) return 0; final Object[] items = this.items; final ReentrantLock lock = this.lock; lock.lock(); try { int n = Math.min(maxElements, count); int take = takeIndex; int i = 0; while (i < n) { c.add((E) items[take]); items[take] = null; take = inc(take); i++; } // 更新 count 等 return n; } finally { lock.unlock(); } }

它一次性把最多 n 个元素转入目标集合,全程持锁,性能比单个 poll 循环好得多。注意两个坑:第一,目标集合不能是队列自身,否则会抛 IllegalArgumentException;第二,它一次性转入 n 个元素后 count 直接减 n,消费方如果用这块做“等待队列空了再暂停”的判断,需要理解 drainTo 之后队列可能还有剩余元素,注意返回值代表实际转移数量。

5. 尺度与边界:粒度问题、空元素禁令和 remove 的 O(n) 代价

这一节集中聊“用 ArrayBlockingQueue 时容易吃暗亏”的细节。每个点都是我或身边同事在真实项目里踩过的。

5.1 禁止 null 元素的真正原因

ArrayBlockingQueue 在 put/offer/add 里第一行就是 checkNotNull(e)。为什么不允许 null?一个直接原因是 poll 方法用 null 表示“队列为空”:

public E poll() { return (count == 0) ? null : dequeue(); }

如果允许 null 入队,那 poll 返回 null 到底是“取到一个 null 元素”还是“队列为空”?调用方无法区分。JDK 的阻塞队列实现(ArrayBlockingQueue、LinkedBlockingQueue、PriorityBlockingQueue 等)统一不允许 null,就是为了把 null 当作“空/失败”信号保留下来。

如果你的业务数据里确实会出现 null,入队前先包装成 Optional 或者用一个哨兵对象包裹,别想着绕开这个约束。

5.2 removeAt 的元素搬移细节

removeAt 的场景是删除一个指定下标的元素,比如 remove(Object) 能删除中间某个元素。它有两段逻辑:

if (i == takeIndex) { // 删除队头,直接清空推进 items[takeIndex] = null; takeIndex = inc(takeIndex); } else { // 删除非队头元素,需要把 i 到 takeIndex-1 之间的元素往前挪 for (int n = i; n != takeIndex; n = dec(n)) { E e = items[dec(n)]; items[n] = e; } ... }

最坏情况下,要搬移几乎整个队列的元素,复杂度 O(n)。如果一个队列容量很大、又频繁做中间元素删除,这条路径会非常伤 CPU。线上排查时如果发现内存队列 CPU 飙升,先看看是不是有人在用 remove(Object) 做“按 key 取消任务”。

5.3 中断处理:lockInterruptibly 与 await 的响应

put/take 都声明了 throws InterruptedException,使用阻塞语义时会响应中断。这意味着如果队列一直满/一直空,线程挂在 await 上,别的线程可以随时中断它。设计上有个点容易忽略:中断发生时,count 状态并没有改变,所以被中断的 put/take 不会产生“半入队、半出队”的不一致状态。这一点由锁保证,可以放心重试。

5.4 内存可见性:锁的边界就是内存屏障

很多初学者会问:ArrayBlockingQueue 里的 count 和数组元素不加 volatile,为什么多线程下能看到最新值?答案是所有读写都在同一把 ReentrantLock 的临界区内完成。ReentrantLock 基于 AbstractQueuedSynchronizer(AQS),其 lock/unlock 操作包含 volatile 状态的读写,天然具备 happens-before 语义。所以在 ArrayBlockingQueue 里,不需要额外给 items 或 count 加 volatile。这也是为什么源码里它们只是普通字段。

这一点在阅读源码时非常关键:不要看到普通字段就以为存在可见性问题,要看它们是否被锁保护。

5.5 慎用 contains 和 toArray

  • contains(Object o):需要遍历整个数组,O(n)
  • toArray():会创建一个新数组并复制元素,同样 O(n)

这两个方法在队列容量很大时都不便宜。我曾经遇到过一个业务,每次任务提交前先 contains 一下判断“是否重复”,队列容量 5000,生产者 20 个线程,直接导致 contains 成为热点。最后换上 ConcurrentHashMap 做去重标记才解决。阻塞队列擅长的是边界内的高效入出队,不是查找。

6. 选型与实战:ArrayBlockingQueue、LinkedBlockingQueue 和场景匹配

这大概是实际工作中被问得最多的问题:到底用哪个?我把二者的核心差异和选择标准整理一下。

6.1 一个维度分清单锁与双锁的本质分歧

维度ArrayBlockingQueueLinkedBlockingQueue
底层结构循环数组,固定容量单向链表节点,容量可选(无界/有界)
锁结构单锁 + 两个 ConditiontakeLock + putLock,两把锁
内存分配初始化时一次性分配,稳定每个元素一个 Node,持续分配,GC 有压力
容量限制必须在构造时指定有界或无界;无界时 put 永不阻塞(但可能 OOM)
入队出队并发不能真正并行可取可放,两把锁并行度更高
中间元素删除支持,但 O(n) 搬移支持,链表删除 O(1),但需要遍历查找

单从这个表看,LinkedBlockingQueue 似乎全面占优,但实际选型远没那么简单。因为 ArrayBlockingQueue 在“容量确定、任务生命周期短、内存敏感”的场景下有更好的稳定性和可预测性。

6.2 建议:先看场景再做决定

  • 生产者消费者数量都不多(比如各 1~4 个),流量平稳:ArrayBlockingQueue 足够,内存更省
  • 要求严格的容量上限,避免无界增长引发的 OOM:ArrayBlockingQueue 更直观,容量显式写死
  • 对批量消费(drainTo)有较高依赖:两者都能用,ArrayBlockingQueue 的批量转储由于连续数组更高效
  • 非常高的吞吐需求、大量线程竞争:可以先试试 LinkedBlockingQueue 的双锁并行,压测对比再定
  • 任务可能偶尔积压且不想丢数据:选有界队列的容量要大,并结合特征评估拒绝策略

6.3 一个可复用的 FastFail 写流控方案

实际项目里我经常用 offer + 超时实现“软限流”:

ExecutorService executor = new ThreadPoolExecutor( core, max, keepAlive, unit, new ArrayBlockingQueue<>(queueSize) ); // 提交任务时 if (!executor.getQueue().offer(task, 500, TimeUnit.MILLISECONDS)) { // 入队失败,走降级逻辑:丢弃或返回提示 return handleReject(task); }

这样队列满时,提交线程最多阻塞 500ms,不会无限挂起,也不会像直接 put 那样可能把上游线程全部拖死。对于需要快速失败、保护系统的场景,这比单纯用饱和策略更可控。

7. 我踩过的几个坑与最后的调参心得

最后分享几个我实际调试过的 Case,每一个都对应上面某个机制,希望能帮你少走弯路。

第一个坑是“队列容量设为 1”。当时做实时数据管道,两个线程互相传递消息,我图省事把容量设成 1,结果性能惨不忍睹。后来用 JFR 一看,几乎所有时间都耗在锁竞争和条件唤醒上。把容量调到 16 之后吞吐翻了好几倍。原因是容量太小时,生产者一入队消费者立即被唤醒,但消费者抢到锁之前生产者又想继续入队,造成两个线程争同一把锁的频率急剧上升。

第二个坑是“用 remove(Object) 取消待处理任务”。系统里有一批延迟任务,消费者还没处理前,如果收到撤销信号就 remove。队列容量 2000,每天跑批,某天突然 CPU 飙到 90%,排查后发现是 remove(Object) 触发了大量元素搬移,把原本 O(1) 的操作变成了近似 O(n)。后来改成每个任务维护一个状态标记,消费者处理前先判断状态是否被取消,彻底摆脱了对中间元素删除的依赖。

第三个坑是“把队列当消息中心用”。曾经有人在业务里用 ArrayBlockingQueue 做跨模块的异步通信,但队列只是进程内的结构,不具备持久化、重试、多副本能力。应用重启后积压在队列里的任务全丢了。阻塞队列适合做线程间沟通,不适合做跨服务或跨重启的可靠消息通道,这个边界要拎清。

至于调参心得:先根据自己的业务模型算出“峰值积压量”,再乘以一个 1.5~2 的安全系数作为容量。别拍脑袋定大小。容量定得太小会让生产者频繁阻塞,定得太大又浪费内存、拖慢遍历类操作。真实场景里,ArrayBlockingQueue 的容量和生产者消费者线程数是强耦合的,最好在压测环境多测几组参数,用数据说话。

最后一个很小的技巧,但很实用:队列的 toString 方法会一次性拼接所有元素,生产环境千万别在日志里随手打 queue.toString()。曾经有个同事在告警日志里打印了容量 10000 的队列内容,结果告警系统直接被打爆,拼接一万条字符串的耗时也远超预期。真要观察队列状态,用 size()、remainingCapacity() 这类 O(1) 方法就足够了。

ArrayBlockingQueue 的源码看过一遍之后,你会发现它其实是并发编程最好的“教科书案例”之一:锁、条件、循环数组、通知唤醒、内存可见性,全都浓缩在这一个类里。把它吃透,再回去看线程池的 workQueue 参数、看各种消息中间件的本地缓存层设计,思路会通透很多。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询