Chapter 03

内部机制:broker 到底怎么做到的

上一章给出了可靠投递的三层防护,并把"classic 队列确认后仍可能丢"指向了 quorum——这一章下探一层,拆开帧、Erlang 进程、流控与 Raft,解释这些保证从何而来。

本章你将建立的 schema

  • AMQP 帧的字节结构:channel 多路复用为什么必须靠帧里的 2 字节编号
  • 每个队列是一个 Erlang 进程:单队列的吞吐天花板,以及为什么加核救不了它
  • credit 流控与内存/磁盘告警:背压如何一路传导到 TCP,告警为什么阻塞所有 publisher
  • quorum 队列的 Raft 多副本,以及 Khepri 取代 Mnesia 后元数据存储的分区行为

3.1AMQP 0-9-1 帧:channel 多路复用的根

一条 TCP 上传输的最小单位是帧(frame);每一帧都带 2 字节的 channel 编号,broker 靠交错这些帧、按编号还原,实现一条连接上多路 channel 并行。

为什么需要它

01 章 §1.1给出了结论:一条 TCP、多个 channel、靠编号区分,同一 channel 串行所以不能跨线程。那一层是从 API 角度看的。这一节落到字节:channel 不是一条独立的 socket,而是同一条 TCP 字节流里被打上同一编号的那些帧。理解了帧,才理解为什么 channel 既能多路复用、又强制串行——两个性质来自同一个设计。

底层机制(比文档深一层):一帧 = 1 字节帧类型 + 2 字节 channel 编号 + 4 字节 payload 长度 + payload + 1 字节帧尾。帧尾固定是 0xCE(十进制 206):解析器读完声明的长度后,下一个字节必须正好是 0xCE,否则判定字节流错位,直接断开整条连接——这是协议自带的成帧校验。帧类型只有四种:method(类型 1,一次 RPC,如 Basic.Publish)、content-header(类型 2,消息属性 + body 字节数)、content-body(类型 3,消息体原始字节)、heartbeat(类型 8,心跳)。

一次逻辑 publish 不是一帧,而是一串:先一个 method frame(声明"这是一次 publish,发给哪个 exchange、什么 routing key"),紧接一个 content-header frame(带消息属性和总 body 大小),再接 1 个或多个 content-body frame(消息体)。消息体大于 frame_max 时被切成多个 body frame——frame_max 在 connection.tune 阶段由 client 与 broker 协商,默认约 131 KB。

一次逻辑 publish(同一 channel 上严格有序) method 类型 1 Basic.Publish exchange · key content-header 类型 2 属性 · body 大小 content-body 类型 3 消息体字节 content-body 类型 3 体 > frame_max 才出现 每帧前缀都带同一个 2 字节 channel 编号 → broker 按编号归属到同一会话 帧尾 0xCE(206)校验成帧;不匹配即断连
图 3.1一次 publish 被拆成 method + content-header + 至少一个 content-body 帧,依次铺在同一 channel 上。注意:三类帧共享同一个 channel 编号,broker 据此把它们重组为一条消息;任何一帧被另一线程的帧插进来,重组就错位。

channel 多路复用的实现,就是给每帧打编号再交错:channel 1 的 method 帧、channel 2 的 body 帧、channel 1 的 header 帧,可以在同一条 TCP 上首尾相接地传,broker 读到每帧前缀的编号,把它投递给对应的会话状态机。这正是应用避免"每个任务开一条 TCP 连接"的原因——几百个并发任务复用一条 TCP 上的几百个 channel,握手、文件描述符、内核缓冲都省下来。

代价直接来自交错:同一 channel 上的帧严格按到达顺序处理——一个 channel 是串行的。所以一个 channel 不能被多个线程共享:线程 A 刚写完 method 帧、还没写 header 帧,线程 B 插进来写自己的 body 帧,编号一样、序列被打乱,broker 重组时把 B 的 body 当成 A 的消息体,或读到错误的帧尾直接断连。这把 01 章 §1.1"channel 不能跨线程"的结论落到了字节层面的因果。

想一想

发一条 200 KB 的消息(frame_max 默认约 131 KB),这条消息在 TCP 上由几个帧组成?

展开答案(先停 10 秒再点)

至少四个:1 个 method frame + 1 个 content-header frame + 2 个 content-body frame(200 KB 超过单帧上限,body 被切成两段)。注意 method 和 header 各自独立成帧、不受 frame_max 影响——被切分的只有消息体。若把 frame_max 调大到 256 KB,body 就只需 1 帧,但单帧内存占用随之上升,这是吞吐与内存的权衡。

洞察 · 成帧校验

帧尾那 1 个字节(0xCE)看似冗余,作用是廉价的损坏检测:长度字段说 payload 有 N 字节,解析器跳过 N 字节后必须撞上 0xCE;撞不上,说明长度被写错或字节流被污染,broker 不试图恢复,直接杀连接。AMQP 把"宁可断连也不处理半个坏帧"做成了协议级别的约束——这也是为什么 channel 上的协议错误往往表现为整条连接掉线,而不是单个 channel 报错。

3.2每队列一个 Erlang 进程:单队列的吞吐天花板

RabbitMQ 用 Erlang 写成,每个队列、每个 connection、每个 channel 都是一个独立的 Erlang 进程;一个队列的全部工作串行穿过它那一个进程邮箱,被钉在大约一个 scheduler/核上——单队列因此有一个换大机器也抬不动的吞吐上限。

为什么是 Erlang

RabbitMQ 跑在 BEAM(Erlang 虚拟机)上。BEAM 的进程不是操作系统线程,而是极轻量的用户态进程:创建一个只要微秒级、几百字节内存,一台机器可以同时跑几百万个。配上 supervision tree(监督树),一个进程崩溃只重启它自己,不波及 broker 整体。把"每个队列、每个连接、每个 channel 都做成一个进程"在别的语言里是奢侈,在 BEAM 上是惯用法——隔离性和容错就是这么来的。

底层机制(比文档深一层):一个队列就是一个 gen_server 风格的 Erlang 进程(classic 队列是 rabbit_amqqueue_process,quorum 队列是 rabbit_fifo 状态机 + Ra 进程)。进程之间不共享内存,只靠消息传递:publish、ack、consumer 注册,全部变成发给该队列进程邮箱(mailbox)的 Erlang 消息,进程从邮箱里一次取一条、顺序处理。BEAM 的 scheduler 通常每个 CPU 核一个,一个进程在某一时刻只能跑在一个 scheduler 上——于是一个队列的全部吞吐被一个进程、约一个核所限定。

这是 RabbitMQ 最具决定性的性能事实:单个队列的处理能力,被单进程、单核的串行邮箱框死。一条队列打满了一个核,换一台核更多、主频更高的机器——那一条队列快不了多少,因为瓶颈不是机器总算力,是单进程串行。扩展的唯一方向是更多队列 / 分片:把负载摊到多个队列进程上,让它们各占一个核并行跑(如 consistent-hash exchange 把消息按 key 散列到 N 个队列,或用 sharding 插件)。绝不是"加一台更大的机器"。

一台 4 核机器 · 一队列 = 一进程 = 钉在一个核 核 1 队列 H 进程 100% · 满载 核 2 队列 A 进程 12% · 空闲 核 3 队列 B 进程 8% · 空闲 核 4 空转 0% · 闲置 高流量 全压在队列 H 一个进程上 加核 / 换更大机器 → 队列 H 仍只跑在一个核上,吞吐不变 解法:把流量分片到多个队列,吃满多个核
图 3.2四个队列进程分散在四个核上,但高流量全打到队列 H:它的单进程把核 1 打满,核 2/3/4 闲着也帮不上。注意:换更大的机器只会让其他核更闲,队列 H 的天花板纹丝不动——必须按 key 分片成多个队列才能横向吃满 CPU。

场景走查:一个订单系统把全部订单事件发进单个 quorum 队列 orders。压测发现吞吐卡在某个数字,CPU 总利用率只有 30%,但 rabbitmqctl 显示 orders 队列进程的 reductions(Erlang 的 CPU 计量)持续打满。换 16 核机器,吞吐几乎不变。根因:单队列单进程的串行天花板。改法是用 consistent-hash exchange 按 order_id 把消息散到 8 个队列 orders.0..7,每个队列各占一个核并行消费——总吞吐随队列数近线性上升。

分片:用 consistent-hash exchange 把单热队列摊成 N 个 bash
# 启用一致性哈希 exchange 插件
rabbitmq-plugins enable rabbitmq_consistent_hash_exchange

# 声明一个 x-consistent-hash 类型的 exchange,按 routing key 哈希分流
rabbitmqadmin declare exchange name=orders type=x-consistent-hash

# 绑定 8 个队列,weight 写在 binding key 上(这里各为 1)
# 消息按 order_id 哈希落到其中一个队列,8 个队列进程吃 8 个核
for i in $(seq 0 7); do
  rabbitmqadmin declare queue name=orders.$i queue_type=quorum
  rabbitmqadmin declare binding source=orders destination=orders.$i routing_key=1
done
反直觉 · 加核救不了单队列

"吞吐不够就升配机器"在 RabbitMQ 单队列上无效。瓶颈是一个 Erlang 进程的串行邮箱,不是机器总算力。更多核、更高主频,只让其他核更闲。能横向扩的只有队列数量:分片、多队列、多消费者。设计阶段就要问"这条队列的峰值会不会顶满一个核",而不是上线后再加机器——加机器这条路对单队列是堵死的。这一点和 Kafka 靠 partition 数扩并行是同一思路,只是 RabbitMQ 的并行单元是"队列"。

想一想

给单个队列加更多消费者,能突破它的吞吐天花板吗?

展开答案(先停 10 秒再点)

只能缓解消费侧,突破不了队列本身的天花板。消息的路由、入队、出队调度、状态记账全在那一个队列进程里串行完成——加消费者只是让出队后的业务处理并行,队列进程本身仍是单核串行的瓶颈。生产侧(publish 进同一队列)更是完全不受益。真正横向扩的唯一办法是把消息分片到多个队列进程,让队列这一层本身并行起来;消费者数量是第二位的。

3.3credit 流控与内存/磁盘告警:背压如何传导

队列给它的发布 channel 发放 credit(额度),每条消息扣 1,队列把消息推向下游(持久化)后才补额度;额度耗尽 → channel 进程阻塞 → TCP 不再读 socket → 背压沿 TCP 直接顶住生产者。

为什么需要它

生产快、消费慢时,消息会在队列里堆积,最终吃光内存。最坏的做法是无限收下再 OOM 崩溃。RabbitMQ 用 credit-based flow control(基于额度的流控)做精准背压:哪条发布路径压垮了哪个队列,就只减速那条路径,不殃及无辜。这比"全局限速"或"丢消息"都更可控。

底层机制(比文档深一层):每个 channel 进程向它要写入的队列进程持有一笔 credit。每发一条消息扣 1 个 credit。队列进程只有在把这条消息向下游推进之后(比如交给持久化层、写入 message store)才给该 channel 回补 credit。生产者快于队列处理速度时,credit 越扣越少,归零那一刻 channel 进程停止处理新的入站帧 → 它不再从 TCP 连接读取字节 → socket 接收缓冲填满 → TCP 滑动窗口归零 → 生产者的 send 阻塞。背压就这样从队列一路逆向传到生产者的网络栈,不需要应用层协议参与,靠的是 TCP 自身的流控。

这套机制是有靶向的:credit 是 channel ↔ 队列这一对之间的,只有压垮了某队列的那条发布路径会被减速。别的 channel、别的队列照常跑。这与下面要说的"告警"形成鲜明对比——告警是全局的核弹。

内存高水位与磁盘告警是最后一道防线,不是日常流控。两条全局红线:

表 3.1 · 两类背压机制的作用域
机制触发条件作用域 / 后果
credit 流控某队列处理慢于其发布 channel靶向:只阻塞那一条 channel 的 TCP 读,其他连接无感
内存高水位告警broker 内存用量超过 vm_memory_high_watermark(默认 RAM 的 0.6,即约 60%)全局:阻塞集群内所有 publishing 连接,consumer 不受影响
磁盘告警剩余磁盘低于 disk_free_limit(默认 50 MB)全局:同上,阻塞所有 publisher,直到磁盘回到阈值以上

底层机制(告警这一层):内存或磁盘告警一旦触发,broker 给所有携带 publish 的连接下达 channel.flow / 直接停止读取,整个集群范围内的生产全部冻结,直到资源回落到阈值以下自动解除。关键的两点:① 告警只冻结 publisher,不动 consumer——让消费继续排空队列,是它能自愈的前提;② 告警只阻塞、不丢消息、也不把消息换页丢弃。它是"踩了急刹",不是"扔货减重"。这与某些系统"内存满了就丢老消息"的策略截然不同。

想一想

broker 内存触顶、所有 publisher 被告警冻结。此时消费者还能正常消费吗?这个设计为什么是对的?

展开答案(先停 10 秒再点)

能。告警只阻塞 publisher,consumer 照常拉取并 ack。这是刻意的:内存压力来自堆积的消息,只有让消费继续、把队列排空,内存才会回落、告警才能解除。若把 consumer 也冻结,系统就死锁在高水位再也下不来。所以"冻结写、放行读"是自愈的唯一可行方向。副作用:生产端会经历一段 publish 卡住或 confirm 迟迟不返回——这正是 02 章 §2.1 publisher confirms 在告警期间"长时间不 ack"的根因,不是消息丢了,是被刹住了。

场景走查:消费者集体宕机,消息在队列里疯涨。先是 credit 流控把对应发布 channel 一条条减速(靶向);堆积继续逼近 60% RAM,内存告警触发,整个集群所有生产连接被冻结,监控里 publish rate 瞬间归零、生产端线程卡在 send 上。运维重启消费者后队列开始排空,内存回落到水位线以下,告警自动解除,publish 恢复。全程没有一条消息丢失,代价是生产侧的一段冻结。

洞察 · 两级背压的层次

把背压看成两级:第一级 credit 流控是外科手术——精确减速闯祸的那条路径,平时就在工作,多数人察觉不到;第二级内存/磁盘告警是全局熔断——第一级没拦住、资源真的逼近极限时,一把冻结所有写入保命。理解这个层次,就能解释一个常见困惑:为什么"只有几个队列在堆积,结果全集群的 publisher 都卡住了"——那是第二级告警被触发了,它从不区分是谁把内存吃满的。回看 01 章 §1.1说的"一条连接被阻塞,其上所有 channel 一起阻塞",告警正是制造这种"一损俱损"的源头。

3.4quorum 队列与 Raft:多副本一致性

每个 quorum 队列本身就是一个 Raft 集群(建在 Ra 库上):一个 leader、N 个 follower、一份复制的预写日志;publisher confirm 只在消息被提交到多数派(N/2 + 1)之后才返回。

为什么需要它

02 章 §2.3指出 classic 队列即使返回了 confirm,落盘也可能在 fsync 之前发生在崩溃窗口里——确认了仍可能丢。要堵上这个窗口,需要的不是"写一份盘写得更勤",而是"把消息复制到多台机器、多数派都记下来才算数"。单机再可靠也扛不住整机宕机;多副本一致性才是。quorum 队列就是 RabbitMQ 对这个问题的答案,用的是工业界验证过的 Raft 共识算法。

底层机制(比文档深一层):每个 quorum 队列是一个独立的 Raft 组,跑在 Ra(RabbitMQ 团队实现的 Raft 库)上。组里一个 leader 负责接收所有写入,N 个 follower 复制 leader 的预写日志(write-ahead log,每条消息是一条 log entry)。一条 publish 的流程:client 发给 leader → leader 追加到自己的 log → 并行复制给 followers → 当多数派(含 leader,共 N/2 + 1 个节点)都把这条 entry 写进各自的 log,这条 entry 被标记为 committed → leader 此时才向生产者回 confirm。confirm 的语义因此变成"已被多数派持久记录",这正是 02 章 §2.3所指的那个 fsync/安全保证的来源。

Producer publish Leader 写 log entry Follower 1 复制 + 落盘 Follower 2 复制 + 落盘 replicate replicate confirm 3 副本,多数派 = 2(leader + 1 个 follower)记下日志 → entry committed → confirm 此刻才返回,不等第 3 个
图 3.33 副本 quorum 队列:leader 写入并向两个 follower 复制,只要多数派(2 个节点,leader + 任一 follower)落了日志,entry 即提交,confirm 立即返回。注意:confirm 在多数派提交后就返回,不必等最慢的第三个节点——这既是它比单机更安全的原因,也是它比 classic 队列延迟更高的原因。

故障与恢复:leader 宕机,剩余 follower 通过 Raft 选举在数百毫秒内选出新 leader,队列继续可用——只要多数派节点还活着。一个宕机后重新加入的节点,从它自己 log 的 offset 处续传缺失的 entry,不需要全量重新同步。这套行为是 Raft 算法保证的,不是 RabbitMQ 自己拼的。

这取代了 classic mirrored queues(镜像队列)。老的镜像队列用的是一套自研的链式复制协议——master 把消息异步推给 mirror,mirror 常常处于"未同步"(unsynchronised)状态,不是任何共识算法。网络分区下它的行为是未定义的:可能两边都自认 master,愈合时一边的数据被丢弃,造成真实的消息丢失。RabbitMQ 4.0 已彻底移除镜像队列(见 01 章 §1.5的过时配置警告)。quorum 队列把"多副本"从一个尽力而为的复制,升级成了有数学证明的多数派共识。

表 3.2 · quorum 队列 vs 已移除的镜像队列
维度quorum 队列(Raft)classic mirrored(已于 4.0 移除)
复制协议Raft 共识,多数派提交自研链式复制,异步、常不同步
分区行为定义明确:少数派侧不可写,多数派继续未定义,愈合时可能丢消息
confirm 语义多数派持久记录后才返回取决于 mirror 是否已同步,不确定
恢复按 log offset 增量续传常需全量重新同步
代价 · quorum 不是免费的

多数派复制要等网络往返,publish 延迟高于 classic 队列;每条消息在内存里维护一份索引(约 32 字节/条),深 backlog 下内存随消息数增长;quorum 队列总是持久的,没有"非持久 quorum 队列"这种东西。换来的是不丢消息与明确的分区语义。选型口径:要可靠性用 quorum(4.x 默认),明确可丢、要极致低延迟的临时流量才用 classic——这条权衡 04 章 §4.2会和 Kafka 的副本机制并排再算一遍。

想一想

一个 5 副本的 quorum 队列,多数派是几个?同时挂掉 2 个节点,队列还能写吗?挂掉 3 个呢?

展开答案(先停 10 秒再点)

多数派 = 5/2 + 1 = 3。挂 2 个,还剩 3 个 = 多数派,能写,confirm 正常返回。挂 3 个,只剩 2 个 < 3,凑不齐多数派,队列停止接受写入(拒绝或阻塞),直到至少恢复到 3 个节点。这正是 Raft 的安全性所在:它宁可在少数派侧停写,也绝不在凑不齐多数的情况下确认——避免脑裂下的数据分叉。代价是可用性:副本越多越能容错,但"多数派"门槛也越高。生产上 quorum 队列通常配 3 或 5 副本,偶数副本不划算(4 副本的多数派也是 3,容错能力却和 3 副本一样只是多耗一份存储)。

3.5Khepri 取代 Mnesia:元数据存储与消息存储

集群元数据(vhost、队列定义、binding、权限)的存储从 Mnesia 换成了基于 Raft 的 Khepri;网络分区下只有多数派分区继续推进,少数派暂停并在愈合后追平——且这个行为刻意不可配置。

两类"存储"别混淆

RabbitMQ 里有两套独立的存储,常被搞混。一套是元数据存储(metadata store):存的是集群拓扑——有哪些 vhost、队列、exchange、binding、用户权限,全集群必须看到一致的视图。另一套是消息存储(message store):存的是消息体本身。前者关乎"集群对自己的认知一不一致",后者关乎"消息字节存在哪"。这一节先讲元数据存储的世代更替,再讲消息存储的落盘细节。

元数据存储:Mnesia → Khepri 的世代更替

底层机制(比文档深一层):Mnesia 是 Erlang 自带的分布式数据库,RabbitMQ 从诞生起一直用它存元数据。它的复制是点对点、无自动冲突解决的:网络分区时,两侧节点都还能各自写元数据(建队列、改 binding),愈合那一刻两份不一致的状态撞在一起,没有自动裁决谁对。官方给的两个缓解手段都不可靠——pause_minority(少数派暂停)在抖动网络里容易误判,autoheal(自动愈合)会挑一侧丢弃;而 pause_minority 触发恢复时甚至可能把整个元数据库 dump 出来再重灌一遍,代价高昂。

Khepri 是 RabbitMQ 团队新写的元数据存储,同样建在 Ra(即 quorum 队列用的那个 Raft 库)之上——元数据的每一次变更都走 Raft,提交到多数派才生效。分区时的行为因此变得确定:只有多数派那一侧能继续修改元数据,少数派一侧暂停元数据写入,等网络愈合后从 Raft 日志追平。和 pause_minority 不同,这不是一个可调的策略,而是 Raft 内建的、唯一的行为。

表 3.3 · Khepri 取代 Mnesia 的时间线(RabbitMQ 4.x)
版本时间元数据存储状态
4.02024Mnesia 仍为默认,Khepri 可选(experimental → 稳定演进中)
4.22025-10Khepri 成为新节点默认元数据存储
4.32026-04Khepri 成为唯一存储,Mnesia 被完全移除
洞察 · "少旋钮"本身就是特性

Khepri 的分区行为故意不可配置,这不是功能缺失,是设计取向。Mnesia 时代的 pause_minority / autoheal / ignore 三选一,把一个本该由共识算法保证的正确性问题,推给运维去赌——选错就丢数据。Khepri 直接用 Raft 的多数派规则消灭这个选择:少数派必停、多数派必进、愈合必追平,无旋钮可调也就无从配错。"fewer knobs"在这里是把"可能配错的自由"换成"默认就对的确定性"。对从 Kafka 来的工程师,这等价于把 ZooKeeper/KRaft 那种"元数据本来就该走共识"的直觉,落实到了 RabbitMQ 的集群层。

消息存储:4 KB 阈值与那道性能悬崖

底层机制(比文档深一层):消息体存哪,取决于它和 queue_index_embed_msgs_below(默认 4 KB)的大小关系。大于该阈值的消息,写进每 vhost 共享的 message store——一组引用计数的段文件(segment files),多个队列引用同一份消息体时靠 ref-count 共享,所有引用都 ack 后才回收。小于阈值的消息,直接内嵌进每队列的 queue index(队列索引)。queue index 本身记录每条消息的位置 + 投递/ack 状态,是队列的账本。

这条 4 KB 线是一道真实的性能悬崖:小于 4 KB 的消息内嵌在索引里,无论 backlog 堆多深,单条消息占用的队列内存基本是常数——索引结构紧凑、可批量顺序刷盘。一旦消息超过 4 KB 走 message store,每条消息多一次间接引用、段文件的随机访问与回收开销随之而来。所以"把消息控制在 4 KB 以内"对深队列场景是实打实的优化,跨过这条线性能特征会变。(顺带:classic 队列那个"confirm 前未必 fsync"的窗口,02 章 §2.3所说,物理上就发生在这套 message store / queue index 的落盘时序里。)

查看与调整内嵌阈值(谨慎,影响内存与吞吐特征) bash
# 当前生效值(默认 4096 字节)
rabbitmqctl eval 'application:get_env(rabbit, queue_index_embed_msgs_below).'
# => {ok, 4096}

# rabbitmq.conf 里调整:把 8 KB 以下的消息也内嵌进索引
# 收益:更多小消息走常数内存路径;代价:索引变大、单段刷盘更重
queue_index_embed_msgs_below = 8192
想一想

同一个队列里堆了一百万条 1 KB 的消息和一百万条 8 KB 的消息,两者对队列内存的占用模式有何本质不同?

展开答案(先停 10 秒再点)

1 KB 的消息小于 4 KB,内嵌在 queue index 里——它们的内容随索引一起,单条占用近似常数,深 backlog 下队列内存增长平缓、可顺序刷盘。8 KB 的消息大于 4 KB,进共享 message store,索引里只留一个引用,消息体在段文件中,多一层间接、回收靠 ref-count,随机访问更多。本质区别:4 KB 线两侧是"内嵌常数开销"与"外置引用 + 段文件管理"两套不同的内存与 IO 模型。批量发小消息时把单条压到 4 KB 以内,能稳稳吃到内嵌路径的红利。

§本章 self-check

先合上教程,把答案写在纸上或编辑器里。写完再点开对照——直接点开等于把这一节当再读一遍。

  1. 一次逻辑 publish 在 TCP 上由哪几类帧组成?帧尾那 1 个字节(0xCE)有什么用?
  2. 一条队列把一台机器的某个核打满了,吞吐还不够。换一台核更多、主频更高的机器有用吗?为什么?正确的扩展方向是什么?
  3. 内存高水位告警触发时,broker 对 publisher 和 consumer 分别做什么?为什么要这样区别对待?它会丢消息吗?
  4. (设计级)一个对延迟敏感、又要求绝对不丢的支付流水队列,应当怎样在 quorum 队列的"多数派复制"代价和可靠性之间取舍?副本数选 3 还是 5,依据是什么?
答案(先做完再展开)
  1. 三类依次:一个 method frame(类型 1,声明 publish 的 exchange 和 routing key)+ 一个 content-header frame(类型 2,消息属性 + body 总大小)+ 一个或多个 content-body frame(类型 3,消息体;超过 frame_max 才会有多个)。0xCE 是帧尾标记,解析器读完声明长度后必须撞上它,否则判定成帧损坏、直接断开整条连接——廉价的损坏检测。
  2. 基本没用。瓶颈是单个队列 = 单个 Erlang 进程 = 串行邮箱,约钉在一个 scheduler/核上;加核、提主频只让其他核更闲,那一条队列的天花板不动。正确方向是分片成多个队列(如 consistent-hash exchange 按 key 散列到 N 个队列),让多个队列进程各占一个核并行,总吞吐随队列数近线性扩展。
  3. 告警阻塞所有 publisher(集群范围),放行所有 consumer。这样区别对待是因为内存压力来自堆积的消息,只有让消费继续排空队列、内存回落,告警才能解除;若冻结消费就会死锁在高水位。它只阻塞、不丢消息、不换页丢弃——是急刹不是扔货。
  4. 设计级:可靠性优先选 quorum 队列(4.x 默认就是),接受多数派复制带来的延迟上浮——支付流水"不丢"的权重高于个位数毫秒的延迟。副本数选 3 起步:多数派 = 2,可容 1 个节点宕机,延迟代价适中。要求更高容错(容 2 个节点同时挂)才上 5(多数派 = 3),代价是每条消息多复制两份、写延迟更高。不选偶数:4 副本多数派也是 3,容错力与 3 副本相同却多耗一份存储与带宽。再叠加 02 章的 publisher confirms,确保 confirm 返回 = 多数派已持久记录,才算端到端不丢。
进阶挑战 · 刚好够不着

诊断一个"全集群 publisher 卡死"的现场

线上告警:所有生产服务的 publish 调用集体卡住、confirm 迟迟不返回,但消费者日志显示仍在正常消费、队列深度在缓慢下降。CPU 总利用率不高,只有少数几个队列在堆积。请按本章机制,分两级说明这是哪种背压、根因链条是什么、为什么消费者不受影响、运维该先看哪个指标来确认。

提示(卡住再展开)

这是第二级背压——内存(或磁盘)高水位告警,不是第一级 credit 流控(流控只会靶向减速个别 channel,不会让全集群一起卡)。链条:少数几个队列堆积 → broker 总内存逼近 vm_memory_high_watermark(约 60% RAM)→ 触发全局告警 → 所有 publishing 连接被冻结,于是连那些没在堆积的服务也 publish 卡住、confirm 不返回。消费者放行,所以队列深度仍在降。先看的指标:broker 的内存告警状态(rabbitmqctl status 里的 alarms / memory,或管理界面顶部的红色告警条)。确认是告警后,处置方向是加速消费排空、或临时调高水位争取时间,而不是去重启 publisher——publisher 没坏,是被刹住了。延伸想:为什么"少数队列堆积却拖垮全集群写入"在 Kafka 上不会以同样方式发生?