ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

队列从循环队列到阻塞队列再到消息队列:一文理清容量、边界与重复消费

队列从循环队列到阻塞队列再到消息队列:一文理清容量、边界与重复消费 队列这玩意儿基础到数据结构课本的第三章可工作里真正把它问透的人真不多。前几天群里还有人问我线程池里那个任务队列到底应该选有界的还是无界的我反问他你知道队列满的时候会发生什么吗他说那就抛异常呗。我又问那阻塞队列呢阻塞在哪儿为什么不直接拒绝他愣了一下。这个问题其实问的不是一个函数而是你对队列这个抽象到底理解到什么程度。这篇东西我想直接聊透队列从数组怎么实现一个不浪费空间的循环队列到链表队列的工程边界再到Python的queue.Queue和Java线程池里阻塞队列的选择逻辑最后落到消息队列的重复消费这种生产事故级的话题。适合正在补数据结构基本功的人也适合写业务代码时遇到过消息堆积、重复消费、线程池任务丢了的读者。1. 队列是什么从排队叫号到系统解耦的通用模型1.1 先进先出一个不讲道理的约束队列最核心的一句话就是先进先出FIFO。你用append往队尾塞用popleft或者pop(0)从队头取谁先进来谁先被处理。这个约束听上去简单但它的价值恰恰在于约束本身。生活里最典型的例子是食堂打饭。先来的人先打到菜后到的排后面秩序就这么维持住了。如果打饭窗口允许后来的插队先打虽然每个人单次处理时间没变但队伍的整体公平性会崩等待时间也没法预估。计算机里到处是队列。打印机的任务队列所有文档排队输出Windows里窗口程序的消息队列鼠标键盘事件不会直接调你的函数而是先压进系统队列窗口过程在循环里一条条取出来处理。你要是把一个消息处理函数写得太长整个界面就卡成未响应因为排在前面的消息全被堵住了后面的进不来。操作系统里的进程就绪队列也是一个队列CPU在时间片轮转时按顺序调度。所以队列解决的不只是存数据的问题还有三个更上层的能力顺序保证、缓冲削峰、解耦生产者和消费者。生产者往队列里丢数据后不用等消费者处理完消费者空闲时再从队列里取两边不需要知道对方的存在节奏也可以不一致。这是所有消息中间件存在的基本逻辑也是为什么基础数据结构里的队列实现后面能一路延伸到分布式消息系统。1.2 队列和栈同样都是线性表命运截然相反很多人一开始学的时候会把栈和队列搞混其实对比着理解特别容易。栈是后进先出LIFO。先进去的反而最后出来就像一摞盘子你只能从顶上拿。递归转循环、括号匹配、浏览器的后退功能全是栈的应用。队列则正好相反先进先出先进来先处理。广度优先搜索BFS用队列是因为你得一层一层往外扩散先遇到先访问。DFS用栈因为要沿着一条路走到黑再回溯。维度栈队列出入位置只在栈顶操作队尾入队头出顺序后进先出先进先出典型场景括号匹配、递归改写、撤销操作BFS、任务调度、消息缓冲实现难点较小单指针搞定数组实现要处理环形下标和满判定面试里经常出现一道题用两个栈实现队列。思路是拿一个栈专门接收入队要出队时把元素全部倒进另一个栈再出那个栈的顶。第一次倒的时候元素顺序会反过来第二次再倒就正过来了。这道题考的本质是对两种结构反序能力的理解。反过来用两个队列实现栈也是同类的变形题你只需要保证每次取元素时从队尾逻辑上回退就行。另外很多算法题里的队列安排队列模拟题目表面上花样很多核心就是维护一条序列支持在头部、中部插入或删除考察的是你对顺序和指针的掌控力。这类题用数组循环队列未必方便反而链表队列或者平衡树的思路更顺手。2. 数组实现与循环队列为什么满队列判断是灵魂2.1 朴素数组实现为什么走不远用数组实现队列最直观的做法是维护两个下标front指向队头rear指向队尾的下一个空位。入队时在rear下标写入元素然后加一出队时返回front下标的元素然后加一。看起来每一步都是O(1)但很快问题就来了。假设数组长度为5你入队5个元素后rear 5此时即便前面出队了几个元素、空出了一些位置只要rear到了数组末尾就无法再入队了。而数组前半段明明是空的。这个现象叫假溢出。处理假溢出的朴素方案是整体搬移出队到一定程度把剩余元素整体前移。但这样每次搬移都是O(n)数据量大一点性能完全不能看。队列要作为底层组件必须做到每个操作严格O(1)所以循环数组几乎是唯一答案。还有一个方案是每次出队时把后面的元素往前移动一位这样front永远指向0但入队和出队全是O(n)更不可取。真用到生产环境的队列不会有人这么写。2.2 循环队列的核心环形下标与空满判定循环数组的思想很朴素逻辑上把数组首尾相接rear到末尾后再回到0。这段接缝靠的是取模运算。下标移动公式是rear (rear 1) % capacity front (front 1) % capacity比如容量为5当前rear 4再入队一个元素时rear (4 1) % 5 0下一次入队就写到0号位。这样假溢出问题没了但紧接着冒出一个更经典的问题空和满怎么区分因为循环之后空队列时front rear满队列时也会出现front rear单靠这两个下标根本分不清。业界主流有三种解法。第一种是浪费一个元素位。数组容量为N实际只放N-1个元素队满条件变为(rear 1) % N front。空条件仍然是front rear。这个方案的好处是不用额外字段代价是浪费掉一个槽位。第二种是加一个size计数器。每次入队size出队size--。队空条件是size 0队满条件是size capacity。这种写法更直观也不浪费空间代码也更不容易错。第三种是用标记位flag入队时置为1出队时置为0再结合front rear判断。实际工程里用得少因为逻辑绕容易在并发场景出问题这里不推荐。这里的空满判定是循环队列的命门。判断错了会出现两类事故队满误判导致本来能存的数据被拒绝入队或者队空误判导致出队取出未初始化的脏数据。没写过真实循环队列的人很难体会这种边界bug的恶心程度。比如容量为4的循环队列满的时候front和rear相邻但绝不相等很多人写着写着就把这个关系搞混了。2.3 两种可落地的代码实现对比先看浪费格子的版本。我习惯让构造时传入的容量是实际可用容量内部数组多申请一个位置。class CircularQueue: def __init__(self, capacity: int): self.capacity capacity 1 # 多出一个位置用于区分空和满 self.data [None] * self.capacity self.front 0 self.rear 0 def enqueue(self, value): if (self.rear 1) % self.capacity self.front: raise OverflowError(queue full) self.data[self.rear] value self.rear (self.rear 1) % self.capacity def dequeue(self): if self.rear self.front: raise IndexError(queue empty) value self.data[self.front] self.front (self.front 1) % self.capacity return value再来看带size字段的版本。这个版本在初始化时可以省去“capacity1”的骚操作逻辑也更贴近人的直觉。class CircularQueueWithSize: def __init__(self, capacity: int): self.capacity capacity self.data [None] * capacity self.head 0 self.size 0 def enqueue(self, value): if self.size self.capacity: raise OverflowError(queue full) tail (self.head self.size) % self.capacity self.data[tail] value self.size 1 def dequeue(self): if self.size 0: raise IndexError(queue empty) value self.data[self.head] self.head (self.head 1) % self.capacity self.size - 1 return value def __len__(self): return self.size这两种实现里我更推荐带size的。有个热词叫“以rear和length分别指示环形队列中的队头和队尾”其实就是在说这种实现思路。rear负责指向队尾length负责记录当前元素个数队头和队尾的关系由长度间接推导。比起浪费一格的写法这种思路在维护长度、判断满和空时都更加直白也不容易在front rear这件事上绕晕。用Go写也类似type CircularQueue struct { data []interface{} head int tail int size int cap int } func NewCircularQueue(capacity int) *CircularQueue { return CircularQueue{ data: make([]interface{}, capacity), cap: capacity, } } func (q *CircularQueue) Enqueue(v interface{}) bool { if q.size q.cap { return false } q.data[q.tail] v q.tail (q.tail 1) % q.cap q.size return true } func (q *CircularQueue) Dequeue() (interface{}, bool) { if q.size 0 { return nil, false } v : q.data[q.head] q.head (q.head 1) % q.cap q.size-- return v, true }顺带提一个优化点如果容量的取值固定是2的幂取模运算% capacity可以用位运算 (capacity - 1)替代速度更快。这是很多高性能环形缓冲区的标配做法但前提是容量必须是2的幂否则不能用这个优化写的时候要小心。2.4 长度计算与扩容复制顺序极容易错队列当前长度的计算有两种方式。用size字段的直接返回len(self)就完了。没有size字段时用公式length (rear - front capacity) % capacity这个公式很多人第一次看会懵。它本质上是处理rear已经绕回0front还没绕回的情况加一个capacity再取模负数就被纠正过来了。同理你在扩容时会发现一个陷阱循环队列里元素不是从0号位开始线性排列的而是从front开始绕了一圈。扩容时正确的做法是新建一个容量翻倍的数组。从旧数组的front开始逐个取出元素按顺序放入新数组的0, 1, 2, ...位置。把新的front置为0新的rear置为size。如果直接从下标0开始复制你会发现元素顺序是乱的因为真实的队头在front不是0。def resize(self): old_data self.data old_cap self.capacity new_cap old_cap * 2 new_data [None] * new_cap for i in range(self.size): new_data[i] old_data[(self.head i) % old_cap] self.data new_data self.capacity new_cap self.head 0 self.size self.size # size 不变tail可由 headsize 推导这段代码里最关键的一行就是old_data[(self.head i) % old_cap]它保证复制顺序严格跟队列的出入队顺序一致。数组扩容本身是O(n)的但均摊到每次入队就是O(1)这也是动态队列能保持高性能的原因。3. 链表队列与阻塞队列数据结构到并发库的工程落地3.1 单链表双指针去掉容量的限制数组实现的队列有容量上限和扩容成本链表实现队列则完全没有这个问题。链式队列维护两个指针head指向队头节点tail指向队尾节点。入队时在tail后面挂新节点出队时移除head节点。class LinkedQueue: def __init__(self): self.head None self.tail None def enqueue(self, value): node ListNode(value) if self.tail is None: self.head node self.tail node else: self.tail.next node self.tail node def dequeue(self): if self.head is None: raise IndexError(queue empty) value self.head.value self.head self.head.next if self.head is None: self.tail None return value单链表实现有个细节出队时需要判断链表是否变空变空时tail也要置空否则tail会指向一个已经不在链表里的旧节点下一次入队时会把新节点挂到一个孤儿节点后面。链式队列和循环数组的取舍我在实际项目里这样判断维度定长循环数组动态数组队列链表队列容量固定需扩容可自动扩容无界内存局部性好Cache友好好差节点分散出队入队复杂度O(1)O(1)均摊O(1)实现复杂度高边界判断多中低适用场景嵌入式、定长缓冲通用编程峰值流量不明、任务量动态嵌入式领域的环形缓冲区基本都用定长数组因为内存有限而且不能频繁分配。普通业务代码里如果你能预估上限优先用数组如果流量不可控链表更稳妥。3.2 阻塞队列与Python的Queue别把卡住当Bug从基础数据结构走到并发编程队列发生了一个重要变化开始有阻塞语义。阻塞队列的意思是当队列满了再put调用线程会停在那里等待直到有元素被取走腾出空间当队列空了再get线程也会停住直到有元素进来。Python里最常用的queue.Queue就是这种阻塞队列。put默认阻塞、get默认阻塞。很多人搜python队列queue不堵塞其实就是把阻塞当成了bug。阻塞本身就是设计目标它是生产者和消费者之间的协作机制让生产者没有空间时自动停下。import queue import threading q queue.Queue(maxsize1) q.put(task1) # 下面这行会阻塞直到另一个线程取走 task1 # q.put(task2)如果你明确需要非阻塞行为用put_nowait和get_nowait。会发现put_nowait在队列满时抛queue.Fullget_nowait在队列空时抛queue.Empty。改一下代码try: q.put_nowait(task2) except queue.Full: print(队列已满任务被拒绝)Python里collections.deque是另一个常见选择append和popleft完全不阻塞空队列popleft会直接抛IndexError。它适合单线程或者用锁管理好的场景多线程高频读写时没有queue.Queue的线程安全保护。有人拿deque和多线程配合使用跑着跑着发现数据错乱根因就是它本身不是为并发设计的。线程池的阻塞队列选择问题在Java生态里更典型。执行器框架提供三种常用队列队列特点风险与适配ArrayBlockingQueue有界固定数组容量可控满了拒绝策略生效适合保护系统LinkedBlockingQueue默认无界容量最大Integer.MAX_VALUE无界时任务堆积可能耗尽内存SynchronousQueue不缓存任务直接交给执行线程适合内部短任务的直接交接新手最容易踩的坑就是默认使用LinkedBlockingQueue的无界模式。任务多了它们不会消失而是全堆在内存里OOM之前没有任何报错提示。我建议生产环境优先用有界队列核显决定容量上限配合CallerRunsPolicy这种拒绝策略让执行不过来的任务回退到调用方线程执行。这样既不会丢任务又不会让内存被无界队列拖垮。3.3 阻塞队列的背压效应背压这个概念在线程池里很重要。有界队列其实就是在制造背压当队列满时生产者被阻塞或者任务被拒绝这个压力会向上游传导让上游降速。链路很长的时候背压会一级一级往上传递最终让源头放慢速度。这不是坏事恰恰是系统自我保护的方式。无界队列看起来不会拒绝任何任务但实际上是把内存容量当成了无限缓冲一旦峰值流量远超处理能力内存会被打爆整个进程崩溃。一个直观的经验如果你不确定容量该设多大先用有界队列加监控观察水位而不是默认无界。等到你看清楚了任务积压规律再决定要不要调整容量。4. 从内存队列到消息队列重复消费问题的排查与根治4.1 数据结构队列和消息队列中间件到底差在哪很多人学完队列后第一次接触RabbitMQ、Kafka这类消息中间件时会觉得这不就是个远程队列吗是但不全是。内存队列解决的是进程内的缓冲和协调问题出队操作等于彻底移除元素谁拿到就是谁的。消息队列解决的是跨进程、跨机器的可靠通信问题它把内存队列里的pop换成了确认消费机制。消费者主动拉取或者订阅消息后中间件并不会立刻删除这条消息而是等消费端显式确认处理成功再删除。这个差异直接导致了重复消费这个高频问题。为什么基础队列里没人讨论重复消费因为dequeue执行完元素就没了不会有第二个人拿到同样的东西。消息队列里由于网络超时、消费者崩溃、消息确认丢失等原因中间件可能把同一条消息投递两次。这个语义在分布式系统里叫至少一次投递at least once。如果严格保证消息不丢失通常要接受可能重复这就是为什么“消息队列重复消费问题”会成为一个热门话题。很多团队最终选择at least once配合消费端幂等而不是追求傻瓜式的exactly once因为百分百精确的一次性语义需要额外的去重机制和性能开销不是默认就能得到的。PHP生态里的Laravel队列、Windows生态里的MSMQ本质上也是把同一个抽象放到不同的承载环境里重复消费的风险始终存在。4.2 一个重复消费问题的真实排查过程假设你收到的现象是同一笔订单在数据库里被插了两条记录或者一条订单状态被更新了两次导致发了两遍通知。第一步先确认是不是生产端重复发送。看生产者的日志有没有重试逻辑被触发。很多框架的消息发送在连接异常时会自动重试如果发送方没有做去重同一个业务请求就会被塞进队列两次。第二步确认消费端ack的时机。如果你的消费者在消息还没处理完就自动返回成功中间件认为这条消息已经被处理掉不会重投。如果你是在处理完业务数据之后再手动ack消息处理中崩溃中间件会认为处理失败继续投递——这其实是正确的行为但你没有幂等保护重复就产生了。第三步检查消费并发数。同一个消费者组里多个实例并发消费两条相同消息同时到两个实例处理没有任何锁或唯一约束的话必然重复。第四步也是最重要的治理手段消费端幂等。一种简单的幂等方案是在数据库表上建唯一键。比如有order_id字段插入前如果主键或唯一索引冲突就说明这条消息已经处理过直接跳过。另一种常用方案是用Redis的SETNX做短期去重import redis r redis.Redis(decode_responsesTrue) result r.set(forder:{order_id}, done, nxTrue, ex3600) if result: # 第一次处理执行业务逻辑 process(order_id) else: # 已处理过直接忽略 return这个方案的要点是nxTrue只有当Key不存在时才能设置成功。第一次处理设置成功第二次再来就返回None自然跳过。加上过期时间可以防止Key无限膨胀。第五步检查broker端的重投策略。有些中间件在消费超时后会按指数退避重投如果重投间隔很短高峰期会形成风暴。把这些参数调到合理范围给消费端留出恢复时间。排查点典型现象解决方案生产者重试同一业务产生多条消息发送前按业务ID去重消费端ack时机处理未完成就ack改为处理完再手动ack消费并发同消息被多实例同时处理配消费端幂等加唯一约束broker重投超时后无限快速重投调大重投间隔或次数限制幂等缺失重复执行业务副作用数据库唯一键或Redis SETNX4.3 消息堆积队列系统里最隐蔽的杀手重复消费还不是消息队列唯一的受害者消息堆积更常见。消费端处理速度跟不上生产速度消息就会在中间件里越压越多消息存活时间长的还会因为过期被丢弃或者进入死信队列。排查堆积问题通常看两个指标队列中积压的消息数量消费者的处理速率。用监控工具查看lag积压量就可以定位堆积发生在哪个消费者组。如果你的消费端是批量拉取模式检查批量大小和并发数如果消费逻辑里有外部调用排查是不是外部依赖变慢导致整体吞吐下降。消息堆积的本质和循环队列队列满非常像。内存队列满时有界队列会阻塞生产者消息队列满时中间件不会一直阻塞上游而是让消息越堆越多。两者的共性在于生产速度永远不能假设消费速度跟得上队列只是缓冲不是无限仓库。处理堆积的基本原则是优先加快消费方不能靠无脑扩容生产者。另外在不同工具里队列权限也值得注意。比如老牌调度系统LSF里有bqueues命令查看队列的权限、用户组和优先级你能看到哪些任务在排队、队列是否被禁用。这种队列已经上升到资源调度层面但根子上的先进先出、排序、优先级这些概念和基础数据结构的队列没有本质区别只是周围多了一圈策略而已。4.4 生产环境的避坑清单最后把几年里真正踩过和见过的坑汇总一下。队列容量不要无界。无论内存队列还是消息队列无界意味着你把整个进程的内存当缓冲出问题的方式是OOM而不是优雅降级。消费者一定要做好幂等。加唯一键、Redis锁或者业务状态判断三选一别只靠中间件保证不重复的幻想。阻塞会被误认为卡死。遇到put长时间不返回先看队列是否已经满、是否有消费者在运行而不是贸然重启进程。循环队列扩容的顺序别搞错。复制时从front开始不是从下标0开始否则元素顺序会乱。不要在生产环境手写队列。成熟库的工程细节远多过你的想象但前提是你自己要知道它的容量、阻塞语义和满时的行为。结尾最后多提一句经验如果让我说一个队列方面最值得花时间练透的点我的答案是循环队列的空满判定和扩容复制顺序。这两个点看起来是几行代码的差别但没有真正自己手写过的人遇到生产环境里的队列边界问题往往连方向感都没有。我在实际项目里的习惯是手写循环队列时用size字段少用浪费一格的方式代码可读性高很多临场也几乎不会错工程上则尽早换成经过验证的组件。数据结构学到后来你会发现复杂系统里的那些精巧机制本质都是在重复队列这门基础课里最朴素的约束与权衡。
返回列表