手写阻塞队列:从底层原理到线程池队列选型全解析

手写阻塞队列:从底层原理到线程池队列选型全解析 上个月调一个上报服务的性能问题发现瓶颈根本不在 IO反而卡在一个自己写的“简易队列”上。生产者线程和消费者线程都在死等CPU 烧到 80%吞吐却一直上不去。后来我花了两晚把阻塞队列从底层重写了一遍问题才真正消失。这篇东西不是教科书复读而是把实现阻塞队列过程中踩过的坑、验证过的细节以及线程池里的阻塞队列选择完整梳理一遍适合正在学并发编程的开发者也适合那些在项目里被无界队列拖垮的运维同学。阻塞队列这个名字听起来唬人本质上就是“一个线程安全、并且支持等待/通知机制”的容器。它解决的生产者消费者问题几乎每个需要异步解耦的系统都会遇到。但很多人在项目里无脑用 JDK 自带的队列出了问题也不知道为什么。我建议你至少从零手写一次有界阻塞队列写完后再去看线程池的队列选型会清晰很多。1. 先搞懂阻塞队列到底在解决什么问题1.1 为什么不能自己加锁然后 while 轮询很多初学者处理生产者消费者问题时第一反应是加一把锁然后无限循环去检查队列状态synchronized (queue) { while (queue.size() MAX) { Thread.sleep(50); } queue.add(item); }这段代码表面能用实际上是个大坑。当队列满了生产者线程会一直循环在 sleep 和检查之间一次检查之间隔了 50 毫秒如果某个消费者恰好在这个空档释放了位置生产者最多要白白等 50 毫秒才能发现如果把 sleep 调小CPU 空转率又会直线上升。前者是延迟后者是浪费。我曾经在一个测试环境里用轮询方式跑 200 个生产者线程线程栈里一抓一大片 RUNNABLE 状态但实际都在做无意义检查CPU 直接被打满业务任务都没时间执行。阻塞队列解决的就是这个矛盾队列满时生产者线程主动挂起消费者取走元素后再精确唤醒生产者队列空时消费者挂起生产者放入元素后再精确唤醒消费者。它把“等待”和“通知”的逻辑封装在队列内部调用方只需要关心 put 和 take不需要自行处理复杂的条件判断。1.2 “阻塞”两字的语义put/take 如何挂起与唤醒先看两个核心动作put 往队尾添加元素如果队列已满线程阻塞直到出现空闲位置take 从队头取出元素如果队列为空线程阻塞直到有新元素入队。关键点在于“满了”和“空了”是两种不同的条件如果用一把锁加一个等待集合去实现很容易出现互相等死的情况。想象一个餐厅等位系统没有位置时客人应该去休息区等服务员叫号时再去前台菜做好的时候厨房要能通知出餐口取餐。这两个等待场景是独立的条件不一样。Java 里 Object.wait/notify 只有单一监视器实现起来比较绕ReentrantLock 搭配多个 Condition 就是为了解决这种“多条件等待”问题。一个 Condition 代表一个等待集合notFull 等待集合装着“因队列满而等待的 put 线程”notEmpty 等待集合装着“因队列空而等待的 take 线程”。生产者入队成功后 signal notEmpty让一个消费者苏醒消费者出队成功后 signal notFull让一个生产者苏醒。这样一来唤醒是精确制导不是一嗓子喊醒所有人。1.3 为什么放着现成的队列不用非要自己实现“JDK 里不是有 ArrayBlockingQueue 和 LinkedBlockingQueue 吗抄一遍有什么意义”我当时的想法很简单自己想改一个带“队列水位统计”和“批量腾挪”功能的队列用现成的类做增强反而要 Hack。更重要的原因是只有从底层实现过一遍遇到线程池的阻塞队列选择、遇到线上莫名其妙的卡死才会有真正的手感。不是为了替代 JDK而是为了在项目里做出“不会出问题”的选型决策。写一遍你会彻底明白为什么 ArrayBlockingQueue 用数组、LinkedBlockingQueue 用链表也明白为什么 put 方法里要用 while 而不是 if为什么 condition 的 signal 比 signalAll 更高效。2. 手写一个有界阻塞队列从接口到实现2.1 接口设计与参数规划我先定义一个最简单但有代表性的接口BoundedBlockingQueueE需要支持 put/take再加 offer/poll 的超时变体方便上层做“等待超时则放弃”的控制。内部存储我选择数组环形缓冲区也就是数组头尾相接插入和删除都通过移动下标完成不涉及元素搬移。环形缓冲区的关键字段有这几个public class BoundedBlockingQueueE { private final Object[] items; private int takeIndex; private int putIndex; private int count; private final ReentrantLock lock; private final Condition notEmpty; private final Condition notFull; public BoundedBlockingQueue(int capacity, boolean fair) { if (capacity 0) { throw new IllegalArgumentException(); } items new Object[capacity]; lock new ReentrantLock(fair); notEmpty lock.newCondition(); notFull lock.newCondition(); } }为什么用环形而不是普通数组加搬移想象一个容量为 8 的数组不断从头部出队头部前面的位置就永远空着最终队列明明还有空间却因为尾部到了末尾而无法插入。环形缓冲区把 putIndex 和 takeIndex 对容量取模实现逻辑上的无限循环这是 ArrayBlockingQueue 的标准做法。如果你选择链表入队只需追加节点出队只需移动头指针看起来更简单但每个节点都要维护对象引用GC 压力更大缓存也不友好在固定容量场景下数组往往更稳。2.2 put 方法等待、入队、精确唤醒先看完整代码再逐行拆解public void put(E e) throws InterruptedException { Objects.requireNonNull(e); lock.lockInterruptibly(); try { while (count items.length) { notFull.await(); } enqueue(e); count; notEmpty.signal(); } finally { lock.unlock(); } } private void enqueue(E e) { items[putIndex] e; putIndex (putIndex 1) % items.length; }第一lock.lockInterruptibly()很关键。它表示线程在等待锁的过程中可以响应中断适合在任务取消或超时控制的场景使用。如果用lock.lock()线程一旦被阻塞在锁上就无法响应中断只能硬等锁释放在线程池关闭时会遇到非常难受的“线程退不出”问题。第二while (count items.length)必须用 while不能用 if。等待线程被唤醒后不能想当然认为条件已经满足。比如两个 put 线程同时被唤醒但队列只有一个空位另一个线程抢先入队当前线程醒来后条件又变成了满。while 循环会再检查一次条件不满足就继续 await这就是条件谓词的标准写法。第三入队后调用notEmpty.signal()而不是notEmpty.signalAll()。为什么 signal 就够了因为每次 put 只会让队列的非空状态从“空”变成“不空”最多只有一个 take 线程能拿到这个新增元素。调用 signalAll 会把所有等待的 take 线程全部唤醒但最终只有一个线程能抢到锁拿到元素其他线程醒来后还要重新进入等待白白增加了上下文切换和锁竞争。这个细节在低并发下看不出来在几百个消费者线程的场景下差别非常明显。2.3 take 方法对称实现与内存泄漏隐患take 的逻辑和 put 完全对称但有一个检查点特别容易踩坑public E take() throws InterruptedException { lock.lockInterruptibly(); try { while (count 0) { notEmpty.await(); } E e dequeue(); count--; notFull.signal(); return e; } finally { lock.unlock(); } } SuppressWarnings(unchecked) private E dequeue() { E e (E) items[takeIndex]; items[takeIndex] null; takeIndex (takeIndex 1) % items.length; return e; }很多人写 dequeue 只取数组元素忘了把items[takeIndex]置空。这在长期运行的队列里会变成诡异的内存泄漏队列逻辑上已经为空但数组引用还死死抓着对象GC 永远回收不掉。尤其是队列里跑的是大对象或者数据库连接池资源时线上 OOM 都可能因此出现。阻塞队列是长期存活的共享对象内部数组的每个空位都应该及时清成 null。take 方法同样在成功出队后 signal notFull每次 take 至多腾出一个空位所以只需要唤醒一个生产者不需要唤醒所有生产者。对称的 while 检查同样不能丢否则队列空了还去取会拿到 null 或者越界。2.4 offer/poll 的非阻塞与超时版本实际项目里put/take 的无限阻塞不够灵活我更常用的是带超时时间的 offer/poll。非阻塞版本逻辑最简单用tryLock拿不到锁就直接返回失败或者拿锁后检查状态不满/不空即刻执行。这里重点讲超时版本因为它藏着一个非常经典的 bug。public boolean offer(E e, long timeout, TimeUnit unit) throws InterruptedException { Objects.requireNonNull(e); long nanos unit.toNanos(timeout); lock.lockInterruptibly(); try { while (count items.length) { if (nanos 0) { return false; } nanos notFull.awaitNanos(nanos); } enqueue(e); count; notEmpty.signal(); return true; } finally { lock.unlock(); } }synchronized时代的Object.wait(timeout)会把剩余等待时间返回而 Condition 的awaitNanos(nanos)同样返回“剩余的纳秒数”。所以必须把返回值重新赋给 nanos再次进入循环判断。如果写成notFull.awaitNanos(nanos);然后不更新变量下一次循环判断的 nanos 还是原始值看起来只是多等了一会儿但极端情况下可能导致无限循环或超时不准。我当时就是在这里踩坑测试用例一直随机挂排查半天才发现是超时时间没有重新赋值。poll 方法是完全对称的就不展开贴代码了核心就是把“不满”换成“不空”把 enqueue 换成 dequeue。3. 边界条件与性能调优把代码打磨到能上线3.1 条件谓词、虚假唤醒与 while 循环我在上一节提过while很重要这里单独展开说说原因。JVM 规范并没有禁止 Condition 的实现产生“虚假唤醒”也就是说即使没有任何线程调用 signalawait 的线程也可能因为底层实现原因醒过来。虽然现代 JVM 上出现概率极低但并发代码必须把它当作必然发生来对待。更实际的原因是即使没有虚假唤醒多个等待线程被 signal 唤醒时也无法保证每个线程醒来后条件都满足。比如队列里有一个空位两个 put 线程都等在 notFull 上一个消费者 take 之后调用了一次 signal两个 put 线程同时被唤醒。第一个抢到锁的入队成功第二个抢到锁时队列又满了如果没有 while 再检查一次它就会带着满队列继续执行 enqueue 覆盖数组数据队列的 count 和实际元素数量就彻底对不上了。所以规范的写法永远是这样while (!conditionIsTrue) { condition.await(); }等不到条件就一直等只有条件为真才继续往下走。这个模式叫“条件谓词循环”是所有基于等待/通知的并发容器都必须遵守的铁律。3.2 中断处理与锁中断策略阻塞队列里的阻塞点基本上是三类等待锁、等待 notFull、等待 notEmpty。等待 Condition 的时候如果线程被中断await 会抛出 InterruptedException 并释放锁等待锁的时候lockInterruptibly也会响应中断。但如果你为了省事用了lock.lock()中断信号就会被挂在一边等线程拿到锁进入 await 之后才被发现线程的退出行为就会延迟导致线程池关闭时“卡住半秒”。还有一个更隐蔽的错误在 try 块里 catch 住 InterruptedException 并默默吞掉。比如这样try { notFull.await(); } catch (InterruptedException ignored) { }这样做的后果是外部线程想中断这个生产者线程线程却对中断视而不见继续闷头执行并返回成功上层的取消机制完全失效。正确的做法是让 InterruptedException 继续向上抛由调用方决定是清理现场还是恢复中断标记。如果实在不能抛至少要调用Thread.currentThread().interrupt()把中断标记恢复回去。3.3 公平锁与非公平锁的取舍ReentrantLock 构造时可以传 fair 参数公平锁保证等待时间最长的线程优先获得锁非公平锁则允许“插队”。ArrayBlockingQueue 的构造方法里也有一个 fair 参数默认是非公平。我刚实现队列时选了公平锁测试时一切正常一到高并发压测性能直接掉了一截。原因是公平锁为了维持 FIFO 顺序在锁竞争激烈时需要做大量的队列检查和唤醒上下文切换非公平锁允许新线程直接抢占虽然可能让少数线程饿肚子但整体吞吐量明显更高。在阻塞队列这个场景我之前也担心“非公平会不会让某个生产者永远插不上队”但实际统计下来除非你极端地让生产者线程数量远多于消费者并且任务量持续饱和否则短时间内不至于真正饿死。真实项目里我更倾向于用非公平锁然后通过合理的队列容量和超时机制兜底。如果一定要保证严格的先来后到比如某些结算任务的生产者必须按序入队那再考虑公平锁但必须接受吞吐损失。3.4 一个实用增强批量操作与水位监控自己手写队列最有价值的一点就是可以按业务场景加能力。我用得最多的是两个增强批量 drainTo 和队列水位统计。批量 drainTo 的好处是减少锁获取次数。比如消费者每次从队列里拿 100 条消息再去批量处理如果每次只 take 一条就要在锁内做 100 次读写切换用 drainTo 一次把队列里所有元素倒到本地链表然后一次性释放锁处理完再继续循环。这个优化在消息批量上报场景里非常明显锁竞争从“每条消息一次”降为“每批一次”。水位统计则更像一个监控埋点。我会在 put 成功后记录当前的 count并更新一个 volatile 的 highWaterMark这样“队列最深到了多少”随时可以查出来用来验证队列容量设置是否合理。如果 highWaterMark 经常逼近容量上限说明消费速度跟不上要考虑加消费者线程或者调整线程池参数。JDK 自带队列也能通过 size 去查但那需要额外加锁而且拿不到历史峰值。4. 常见问题与排查技巧实录4.1 队列空了take 线程却一直卡住这种问题最典型的场景消费者线程全部停在 notEmpty 上队列里明明有数据可它们就是不醒。我第一时间会抓线程栈确认等待位置然后检查 put 成功之后是否真的调用了notEmpty.signal()。这里有个容易写错的点把 notFull 和 notEmpty 搞混入队之后误调用了 notFull.signal()那在场的 take 线程当然没人叫醒。还有一个坑是 signal 写在了锁外面比如先unlock()再signal()。理论上 Condition.signal 并不要求一定在锁内但如果两个线程一先一后操作很容易发生信号丢失一个线程调用 signal 之后准备唤醒 take另一个线程插进来把状态又改回去了就等于叫醒的动作被“覆盖”了。规范做法是在锁内、finally unlock 之前完成 signal确保和状态变更在同一个临界区内。4.2 容量明明没满put 依然一直阻塞反过来地队列还远没到容量上限生产线程却全挂在 notFull 上。八成是 count 和实际数组元素数不一致。最常见的是 dequeue 后忘了把 count--或者 try 块里某个分支提前 return 前忘了维护 count。手动实现阻塞队列时count 就是整个容器正确性的命根子建议每次入队/出队后都打印一次调试日志用单线程跑一轮 put/take确认 count 在边界值处一致。还有种返祖现象是有人在 enqueue 内部把 putIndex 更新了但 count 更新放在了锁外这会导致其他线程进入临界区时看到 inconsistent 的状态。我在很早之前犯过这个错直接把 count 写在 finally 外面结果两个线程同时入队count 只加了一次队列实际有两个元素却显示一个随后 take 线程只拿到一个元素就以为空了。4.3 高并发下吞吐量上不去如果正确性没问题但压测时吞吐一直上不去优先怀疑三件事锁竞争、惊群效应、取模运算开销。锁竞争最直接的解法是缩短临界区。enqueue 和 dequeue 本身很简单不要在里面塞其他耗时逻辑比如打印日志、做统计、触发外部回调这些全部挪到锁外面。惊群发生在误用 signalAll 时我前面已经说过应该用 signal 精确唤醒。取模运算则是一个细节优化如果把队列容量设计成 2 的幂那么(putIndex 1) % items.length可以换成(putIndex 1) (items.length - 1)位运算比取模快一截。但这类优化不是说不能做而是要在 profiling 确认瓶颈后再做不要一开始就把代码写成位运算让可读性变差。我压测时用 jstack 抓线程状态如果看到大量线程处于 TIMED_WAITING 而不是 RUNNABLE说明它们真的在休眠锁竞争可能没想象中严重如果看到大量 RUNNABLE 却都在LockSupport.park附近打转那才是锁竞争热点。4.4 线程池队列选型导致的隐蔽问题这个问题严格来说不是手写队列本身的 bug但很容易误判成队列 bug。ThreadPoolExecutor 的 workQueue 如果选的是无界队列会出现一个反直觉的现象核心线程耗尽之后新任务不会去创建非核心线程而是全部堆到队列里。于是线程池最大线程数形同虚设高峰期的积压任务全在内存里排队等到响应超时内存也快满了。我见过一个线上服务就是用了new LinkedBlockingQueue()默认无界容量半个小时内任务堆积上百万个最终直接 OOM。线程池的阻塞队列选择必须提前算清楚不能靠默认参数混过去。5. 线程池的阻塞队列选择别再什么都用无界队列5.1 ThreadPoolExecutor 参数与队列的联动ThreadPoolExecutor 有七个核心参数其中 workQueue 决定了任务提交后的“缓冲策略”。任务执行流程可以概括成四步核心线程数未满时直接创建新线程执行核心线程满后优先把任务放入 workQueue队列满后才继续创建线程直到达到 maximumPoolSize超过最大线程数后触发拒绝策略。这个流程里队列的容量直接影响“队列满”这个触发点何时到来。队列越大非核心线程越晚被创建任务排队时间越长队列越小线程数能更快扩展但也会更早到达拒绝策略。很多人以为最大线程数设得越高越好但如果用了无界队列最大线程数根本没有上场机会。5.2 常见队列的特性对比队列是否有界特点典型使用场景ArrayBlockingQueue有界数组实现容量固定可配置公平需要严格控制资源占用的线程池LinkedBlockingQueue默认无界可指定容量链表实现吞吐较高无界时可能 OOM有界用法适合大多数异步任务SynchronousQueue不存储元素每个 put 必须等待一个 take适合直接移交希望任务不排队、立即交给线程执行的场景PriorityBlockingQueue无界支持按优先级出队需要优先处理高优先级任务的场景DelayQueue无界元素需要实现 Delayed到时间才能被 take定时任务、延迟队列场景我在项目里最常用的搭配是核心链路低延迟用SynchronousQueue因为任务压根不排队必须就有线程立刻接走普通异步任务用有界LinkedBlockingQueue容量一般设置 1000 到 5000拒绝策略用 CallerRunsPolicy 或自定义的丢弃策略需要延迟执行的任务用固定大小线程池配合DelayQueue每个任务设定延迟时间到期才被消费。5.3 队列容量怎么估算网上很多经验贴只告诉你“用有界队列”却不说容量设多少。我一般用一个小公式估算队列容量 可能同时到达的任务数 - 核心线程数× 期望排队时间 / 单个任务平均耗时举个例子你的系统峰值并发提交 100 个任务核心线程是 8单个任务平均执行 50 毫秒你希望任务最多在队列里等 200 毫秒那么容量大概是(100 - 8) × 200 / 50 ≈ 368所以初始值可以设 400 左右。这个公式不严谨但它有一个好处强制你把“并发数”“执行耗时”“可接受的排队时间”这三个关键数字先想清楚。上线后要配合监控持续调整核心指标是队列的 highWaterMark如果长期接近容量上限就调大如果长期很低就调小。千万别把无界队列作为“简单省事”的默认选择。无界队列表面上永远不会触发拒绝策略实际是把内存风险无限放大了。宁可让拒绝策略介入也不要用内存去硬扛流量洪峰。5.4 实际选型经验与踩坑记录我线上有个上报服务早期用的是无界 LinkedBlockingQueue 固定线程池高峰期积压任务太多消费者线程处理不过来队列在内存里一直膨胀。后来我把队列改成了有界 ArrayBlockingQueue容量 2048拒绝策略用 CallerRunsPolicy。这样做的效果是一旦积压超过 2048新任务不再无限入队而是由提交任务的线程自己去执行。虽然提交线程被拖慢但系统不会 OOM业务也能感知到“再往下压就要排队了”这在很多场景里反而是一种健康的背压机制。SynchronousQueue 则要特别小心。因为队列不缓存任务如果线程池里没有空闲线程提交操作会直接阻塞直到有线程接手。它配合 maximumPoolSize 比较大的线程池可以让非核心线程快速创建但如果 maximumPoolSize 设得太小提交线程会频繁阻塞吞吐量远低于预期。如果只是想要“不排队、马上执行”但又不希望提交线程被阻塞建议用有界队列 较大的最大线程数。6. 写在最后自己实现一次比背十篇面试题都有用6.1 我学到的东西手写阻塞队列这件事最深的体会不是造轮子有多爽而是它把“并发编程到底在并发什么”这个问题彻底打通了。以前看 Condition 的源码总觉得是抽象概念自己实现一遍才发现Condition 的等待集合、await 释放锁、signal 唤醒每一环都必须和业务状态严格对应少一个 signal多一个 if都会造成线上事故。还有就是对性能的理解。我原先以为高并发队列只要加锁就行实际测试下来公平锁和非公平锁、signal 和 signalAll、取模和位运算这些细小的选择在热点路径上会被放大很多倍。所以后来看 ArrayBlockingQueue 的源码不再觉得那是一堆“死代码”而是能看出作者在每个细节上的取舍。6.2 可以继续扩展的方向如果你的目的是学习我建议再往这几个方向试试给队列增加关闭状态关闭后 put 直接拒绝、take 把剩余元素耗尽支持按优先级出队把数组换成二叉堆或者把锁换成 CAS 实现一个无锁队列再去对比性能和复杂度。手写阻塞队列不是项目终点而是理解线程池、理解消息中间件、理解分布式排队系统的一块跳板。最后再分享一个实操小技巧测试阻塞队列时别一上来就开 100 个线程压测。先用单线程把 put/take 的边界测清楚再用 2 个生产者、2 个消费者去验证唤醒逻辑。我习惯在测试里给每个 put/take 打上时间戳出现挂死时能立刻从日志看出是哪个条件没被满足。这样定位问题的速度比用调试器打断点快得多。