资讯动态

并发 10 · 阻塞队列与生产者-消费者

发布时间:2026/8/21 6:43:31 来源:尧图企业网站定制
秒杀系统里有一个经典拆分思路下单请求先快速写入一个队列返回排队中后台线程再慢慢从队列里取出请求做真正的扣库存、生成订单。这样一来前端接口响应飞快即使瞬时涌入十万请求也不会直接怼到数据库上——这正是生产者-消费者模式在真实业务里的样子下单请求是生产者扣库存的后台线程是消费者中间那个队列就是这一篇的主角。第 8 篇讲线程池的时候我们已经见过这个队列的身影——ThreadPoolExecutor内部的任务排队靠的正是BlockingQueue。这一篇要把它专门拎出来讲透它到底解决了什么问题家族里几个常见实现各有什么脾气以及生产者-消费者这个模式怎么落地成代码。最后我们会跳出 JDK 自带工具简单看一眼工业级消息组件 Disruptor 的设计思路作为还能怎么更快的参照。线索是先看不用BlockingQueue、手写wait/notify有多繁琐理解BlockingQueue到底帮我们封装了什么然后拆开看ArrayBlockingQueue、LinkedBlockingQueue、SynchronousQueue三个最常用实现的结构差异再看两个带排序语义的特殊队列——PriorityBlockingQueue和DelayQueue接着用BlockingQueue走一遍秒杀下单排队消费的完整代码最后看看 Disruptor 用什么招数把吞吐再往上顶一层。目录不用 BlockingQueue 会怎样手写 wait/notify 的痛点BlockingQueue四组方法把等待封装起来ArrayBlockingQueue vs LinkedBlockingQueue单锁与双锁分离SynchronousQueue不存东西的队列PriorityBlockingQueue 与 DelayQueue带排序语义的队列实战用 BlockingQueue 搭一个秒杀下单排队消费一瞥 Disruptor无锁环形缓冲区选型小结一、不用 BlockingQueue 会怎样手写 wait/notify 的痛点在有BlockingQueue之前生产者-消费者模式要靠synchronizedwait()/notify()手写。逻辑并不复杂共享一个缓冲区消费者发现缓冲区空就wait()生产者放入数据后notify()唤醒反过来生产者发现缓冲区满就wait()消费者取走数据后notify()。classBuffer{privatefinalQueueOrderqueuenewLinkedList();privatefinalintcapacity100;publicsynchronizedvoidput(Orderorder)throwsInterruptedException{while(queue.size()capacity){wait();// 缓冲区满等待被唤醒——必须用 while 不能用 if}queue.offer(order);notifyAll();// 唤醒可能在等非空的消费者}publicsynchronizedOrdertake()throwsInterruptedException{while(queue.isEmpty()){wait();// 缓冲区空等待被唤醒}Orderorderqueue.poll();notifyAll();// 唤醒可能在等非满的生产者returnorder;}}这段代码能跑但藏着几个容易踩的坑必须用while而不能用if判断条件。线程被notifyAll()唤醒后并不保证条件一定已经满足——可能被唤醒时缓冲区又被别的线程占满了“虚假唤醒”Spurious Wakeup用if会导致醒来后不复查条件直接往下执行出错。这个细节极易被忘记。notify()还是notifyAll()容易选错。只用notify()可能唤醒的恰好是同类线程比如两个消费者互相唤醒而真正在等的生产者永远没被叫到程序卡死。稳妥起见通常用notifyAll()但这意味着每次唤醒都要把所有等待线程叫起来重新抢一遍锁有一定浪费。锁的粒度是整个Buffer对象生产者和消费者共用同一把锁即便逻辑上放和取可以有更细的并发控制这里也做不到。这些细节手写起来很容易出错而且这类缓冲区等待/唤醒的模式在并发编程里太常见了——BlockingQueue就是 JDK 把这一整套逻辑标准化、封装好之后交给我们的现成工具。二、BlockingQueue四组方法把等待封装起来BlockingQueue接口最核心的能力队列满时插入自动阻塞队列空时取出自动阻塞把上一节那堆whilewaitnotifyAll的细节全部封装到了put()/take()内部。它提供了四组语义不同的方法分别应对操作当下做不了该怎么办这个问题处理方式插入移除行为抛异常add(e)remove()队列满/空时直接抛异常不等待返回特殊值offer(e)poll()队列满/空时返回false/null不等待阻塞等待put(e)take()队列满/空时一直阻塞直到条件满足超时等待offer(e, time, unit)poll(time, unit)阻塞等待但超过指定时间就放弃生产代码里最常用的是阻塞等待put/take典型的生产者-消费者场景和超时等待担心一直阻塞卡死线程给个上限超时后走降级逻辑。add/remove那组抛异常的方式很少直接用在业务代码里。内部实现上ArrayBlockingQueue、LinkedBlockingQueue这类实现都是靠ReentrantLockCondition来做等待/唤醒回顾第 6 篇讲的Lock家族——Condition本质上就是wait/notify的等价物只是可以从同一把锁上分裂出多个独立的等待队列比如队列非空和队列非满用两个不同的Condition唤醒时可以精准叫醒该叫的那批线程不用像notifyAll()那样一次性全叫起来再抢。这正是用高层工具替我们封装好细节的价值所在。三、ArrayBlockingQueue vs LinkedBlockingQueue单锁与双锁分离这两个是最常用的两个实现名字看起来只是底层用数组还是链表的区别但背后的锁设计差异更值得关注。ArrayBlockingQueue数组实现有界单锁。创建时必须指定容量之后不能扩容或缩容——这意味着内存占用从一开始就是确定的、恒定的不会因为任务堆积而无限膨胀。内部用一把ReentrantLock配两个ConditionnotEmpty和notFull无论是插入还是取出操作都要先抢这同一把锁。这意味着生产者和消费者在同一时刻互相排斥——即使是生产者在放和消费者在取这两件逻辑上不冲突的事也得靠这一把锁串行化。LinkedBlockingQueue链表实现可选有界双锁分离。不指定容量时默认上限是Integer.MAX_VALUE近似无界这也是第 8 篇讲Executors.newFixedThreadPoolOOM 风险的根源——它内部用的就是没指定容量的LinkedBlockingQueue。关键差异在锁它把放和取用两把独立的锁分开putLock和takeLock生产者之间互相排斥、消费者之间互相排斥但生产者和消费者彼此不冲突可以真正同时进行——这比ArrayBlockingQueue的单锁设计并发度更高。代价是链表节点需要动态创建和回收多了一层 GC 开销性能不如数组那样稳定可预测。// 秒杀场景容量恒定、延迟稳定优先 —— 选数组BlockingQueueOrderorderQueuenewArrayBlockingQueue(1000);// 高并发读写混合、更看重吞吐 —— 选链表务必显式指定容量BlockingQueueOrderorderQueue2newLinkedBlockingQueue(1000);// 别漏了容量参数选型的直觉要稳定延迟和确定内存占用选ArrayBlockingQueue要更高吞吐、能接受一点 GC 波动选LinkedBlockingQueue——但无论选哪个永远显式指定容量这是第 8 篇反复强调的Executors 陷阱的直接延伸。四、SynchronousQueue不存东西的队列SynchronousQueue是这个家族里最反直觉的一个——它内部不存储任何元素,可以理解成容量为 0 的队列。它的语义是直接传递Rendezvous生产者的put()必须一直等到有消费者调用take()来接货两边才能同时完成反过来也一样。这和LinkedBlockingQueue哪怕容量设成 1 都完全不同——容量为 1 的队列允许生产者先把货放进去、转身离开货物在队列里等着什么时候有消费者来取都行SynchronousQueue不允许这种放下就走必须是一次面对面的交接。这个特性用在哪回顾第 8 篇讲过的Executors.newCachedThreadPool()——它的任务队列就是SynchronousQueue。逻辑是任务提交时如果恰好有空闲线程能立刻接手就直接交接执行如果没有空闲线程由于队列不能存任务会立刻触发创建新线程而不是把任务积压在队列里等着。这种不允许任务排队、要么立刻处理要么立刻扩容的设计适合任务处理要求低延迟、不希望任务在队列里干等的场景。五、PriorityBlockingQueue 与 DelayQueue带排序语义的队列前面几种队列都是先进先出FIFO但有两类需求需要按别的顺序出队。PriorityBlockingQueue按优先级出队。底层是一个二叉堆和普通PriorityQueue的排序逻辑一样只是额外包了一层锁让它具备阻塞队列的能力普通PriorityQueue没有take()这种阻塞方法。它逻辑上是无界的出队顺序由元素的自然顺序或自定义Comparator决定而不是插入顺序。典型场景VIP 订单要比普通订单优先处理——把订单按会员等级、下单金额排序塞进这个队列消费者线程永远先拿到优先级最高的那个。DelayQueue只有到期的元素才能被取出。元素必须实现Delayed接口声明自己还有多久到期底层同样基于堆结构但排序依据是剩余延迟时间——离到期最近的排在堆顶。take()会一直阻塞直到堆顶元素真正到期。典型场景订单超时未支付自动取消——下单时把订单包装成一个30 分钟后到期的延迟任务扔进队列专门有一个消费者线程take()阻塞等待一旦有订单到期就被取出来执行检查支付状态、未支付则取消库存释放的逻辑。这比用Timer定时轮询数据库高效得多也比给每个订单单独开一个定时任务线程节省资源。classOrderExpireTaskimplementsDelayed{privatefinallongexpireAt;// 到期时间点毫秒时间戳privatefinalLongorderId;publicOrderExpireTask(LongorderId,longdelayMillis){this.orderIdorderId;this.expireAtSystem.currentTimeMillis()delayMillis;}OverridepubliclonggetDelay(TimeUnitunit){returnunit.convert(expireAt-System.currentTimeMillis(),TimeUnit.MILLISECONDS);}OverridepublicintcompareTo(Delayedo){returnLong.compare(this.expireAt,((OrderExpireTask)o).expireAt);}}// 下单时塞入一个 30 分钟后到期的任务delayQueue.put(newOrderExpireTask(orderId,TimeUnit.MINUTES.toMillis(30)));// 后台专职线程阻塞等待到期订单取到就说明该检查支付状态了while(!Thread.currentThread().isInterrupted()){OrderExpireTasktaskdelayQueue.take();// 未到期的元素不会被取出checkAndCancelIfUnpaid(task.getOrderId());}六、实战用 BlockingQueue 搭一个秒杀下单排队消费回到开头的场景前端下单请求先快速入队返回后台线程池慢慢消费做真正的扣库存。用BlockingQueue搭起来非常直接。publicclassSeckillOrderQueue{privatefinalBlockingQueueSeckillRequestqueuenewLinkedBlockingQueue(5000);privatefinalExecutorServiceconsumersExecutors.newFixedThreadPool(4);// 生产者接口线程调用快速返回publicbooleansubmit(SeckillRequestreq){// offer 而不是 put队列满了直接拒绝而不是阻塞住接口线程returnqueue.offer(req);}// 消费者启动固定数量的后台线程持续消费publicvoidstartConsuming(){for(inti0;i4;i){consumers.execute(()-{while(!Thread.currentThread().isInterrupted()){try{SeckillRequestreqqueue.take();// 阻塞等待队列空就等着deductStockAndCreateOrder(req);}catch(InterruptedExceptione){Thread.currentThread().interrupt();// 回顾第1篇响应中断要恢复中断状态break;}}});}}privatevoiddeductStockAndCreateOrder(SeckillRequestreq){// 真正的扣库存、生成订单逻辑这里会用到第4篇的CAS/AQS锁保证库存扣减的原子性}}几个设计上值得注意的点生产者侧用offer()而不是put()。接口线程是给用户返回响应的绝不能被阻塞卡住——队列满了就直接拒绝返回活动太火爆请稍后重试之类的提示而不是让用户的请求线程也跟着排队等待。消费者侧用take()阻塞等待。后台消费线程本来就是常驻的队列空的时候老老实实阻塞等着就好没必要空转轮询浪费 CPU。队列必须有界这里是 5000。呼应第 8 篇和这一篇第三节反复强调的不设边界的队列在真实的秒杀流量面前就是一颗定时的 OOM 炸弹。这套结构本质上就是用一个有界队列做流量削峰——前端瞬时冲进来的流量被队列缓冲、平滑成后端能匀速处理的节奏这也是BlockingQueue在真实系统里最常见的用法。七、一瞥 Disruptor无锁环形缓冲区BlockingQueue已经能满足绝大多数场景但它终究是基于锁的实现——竞争激烈时线程阻塞、被唤醒都要经过上下文切换的开销回顾第 1 篇讲的上下文切换代价。LMAX一家外汇交易平台在追求极致吞吐单机百万级 TPS时设计了一个更激进的替代方案——Disruptorlog4j2、canal 等知名项目内部都用它来传递消息。Disruptor 的核心结构是RingBuffer环形缓冲区底层是一个固定大小的数组逻辑上首尾相连成一个环生产者不断往环上写、消费者不断从环上读写到数组末尾后又绕回开头继续写。它比BlockingQueue快的关键设计提前分配好整个数组写入时只是填格子没有节点的动态创建和回收——彻底避开了LinkedBlockingQueue那种链表节点带来的 GC 压力。用生产者游标和消费者游标本质是原子变量来协调进度而不是靠锁互斥。生产者知道自己写到了哪个位置、消费者们各自读到了哪个位置通过比较游标位置就能判断这个格子的数据能不能读/能不能写——大部分情况下靠 CAS 更新游标就完成协调不需要线程阻塞等待锁释放。多个消费者可以并行消费同一份数据不同于队列取走就没了的语义也可以通过声明依赖关系确定先后消费顺序比单纯的一个队列配多个消费者互相抢灵活得多。用一句话概括两者的差异BlockingQueue用锁阻塞换来了简单直观的编程模型Disruptor 用无锁的游标协调换来了更高的吞吐但复杂度和学习成本也更高。对绝大多数业务系统而言BlockingQueue的吞吐早已绰绰有余Disruptor 更适合像交易系统这种对延迟和吞吐有极端要求的场景——知道它的存在、明白它解决的是什么问题比在普通业务代码里生搬硬套更重要。八、选型小结需求推荐关键原因容量恒定、延迟稳定ArrayBlockingQueue数组单锁内存占用从创建起就确定高并发读写混合、追求吞吐LinkedBlockingQueue双锁分离生产者消费者互不阻塞任务要立即处理不允许排队积压SynchronousQueue不存储元素强制面对面交接按优先级处理任务PriorityBlockingQueue堆排序VIP/紧急任务优先出队到期才处理超时取消、定时任务DelayQueue堆排序按剩余延迟时间出队极端吞吐/延迟要求金融交易级Disruptor无锁环形缓冲区绕开锁竞争和GC压力一条主线贯穿全篇BlockingQueue把生产者-消费者这个古老的并发模式里最容易出错的等待/唤醒细节封装成了几个方法调用家族里每个实现只是在用什么结构存数据、锁怎么分这两个维度上做了不同取舍。带走两句话① 生产者-消费者场景默认用BlockingQueue别自己手写wait/notify有界还是无界永远要想清楚无界队列是压垮系统的隐形杠杆。② 遇到要排序的需求先别急着自己写比较逻辑排队PriorityBlockingQueue和DelayQueue往往已经是现成答案。下一篇我们进入同步工具类——CountDownLatch、CyclicBarrier、Semaphore、Phaser、Exchanger这几个工具解决的是多个线程之间怎么互相等待、怎么限流协调的问题和这一篇的队列思路互补是并发工具箱里另一组常被面试问到、也确实常用的利器。

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价