Chapter 02
可靠投递:让消息的托管链不断裂
上一章建立了消息从 producer 经 exchange/binding 到 queue、并在 ack 后被删除的路径——这一章把这条托管链的每个接缝补上,让消息不丢。
本章你将建立的 schema
- 发送端的接缝:publisher confirms 用 delivery tag 异步回执,确认"消息已被 broker 托管"
- 消费端的两个接缝:手动 ack 确认"已处理可删除",prefetch 控制"一次推多少"防止撑爆
- 持久化的三个独立开关:durable(实体)≠ persistent(消息)≠ 真正安全(quorum 多副本)
- 失败兜底:dead-letter exchange、TTL 过期、delivery-limit 重试上限,把毒消息引开而非丢弃
2.1发送端不丢:publisher confirms
publisher confirms 让 broker 在消息被处理后回一个确认,生产者据此知道这条消息已被托管,不是发进了黑洞。
裸的 basic.publish 是单向的——TCP write 返回成功只说明字节进了内核缓冲,不说明 broker 收到、更不说明它路由进了队列或写了盘。上一章 1.2 已经埋了这个缺口:没有 binding 命中的消息被静默丢弃,生产者一无所知。publisher confirms 把这条单行道改成有回执:broker 处理完一条消息,回送一个带 delivery tag 的 basic.ack,生产者凭它对账。
底层机制(比文档深一层):confirm 是异步且按 delivery tag 累进的。生产者在 channel 上开启 confirm 模式后,broker 给这条 channel 上发出的每条消息按发布顺序编号(从 1 递增的 delivery tag)。生产者不必发一条等一条——它持续 publish,broker 在消息"处理完成"时回 basic.ack,且可以批量确认:一个 basic.ack 带 multiple=true 表示"到这个 tag 为止全部确认"。所谓"处理完成"对持久化消息意味着"已写入"(写入的真实强度见 2.3,那里有一个关键的 fsync 缺口);路由失败的消息走 basic.nack。生产者侧需维护一个"已发出未确认"的窗口,收到 ack 后从窗口移除,超时或 nack 则重发。
AMQP 还有一套老的事务机制 tx(tx.select/tx.commit),同步阻塞:每次提交都要等 broker 完整落地再返回。官方实测它比 publisher confirms 慢约 250 倍。结论很硬:发送端可靠性用 confirms,不要用 tx。看到代码里还在 txSelect(),基本是十年前的写法。
场景走查:一个支付回调服务每秒 publish 几千条结算消息。开启 confirm 模式后,它把每条消息的 delivery tag 和业务 ID 存进一个本地未决表,继续不停 publish。broker 陆续回 basic.ack(tag=4200, multiple=true),服务一次性把 tag ≤ 4200 的记录从未决表清掉。某条消息的 routing key 拼错、无队列可达,broker 回 basic.nack(tag=4201),服务据此告警并把这条转人工补偿——这条本来会被静默吞掉,现在看得见了。
channel.confirm_delivery() # 开启 confirm 模式
def on_confirm(tag, multiple, nacked):
# multiple=True 表示 tag 及之前的都确认;nacked=True 表示这是 nack
if nacked:
resend(pending.pop_le(tag)) # broker 拒收,重发
else:
pending.discard_le(tag) # 已托管,从未决窗口移除
channel.add_on_confirm_callback(on_confirm)
for msg in stream: # 持续发布,不逐条阻塞等待
channel.basic_publish("orders", msg.key, msg.body,
mandatory=True) # 无队列可达就退回,不静默丢
pending.add(channel.next_delivery_tag, msg)
上面的 mandatory=True 是 confirms 的搭档:confirms 回答"broker 收下了吗",mandatory 回答"收下后有队列接吗"——无队列可路由时 broker 通过 basic.return 把消息退回生产者,而不是默默丢弃。两者合起来才覆盖 1.2 里那个"publish 成功 ≠ 进了队列"的缺口。
生产者收到了 broker 的 basic.ack confirm,能不能就此断定"这条消息一定不会丢了"?
展开答案(先停 10 秒再点)
不能,要看队列类型。confirm 只保证 broker 已按它对这条消息的承诺处理完。对 classic 队列,承诺仅是"写进了内存写缓冲并已下发写盘指令",并未 fsync——broker 此刻崩溃,这条已 confirm 的持久化消息仍可能丢(缺口见 2.3)。只有 quorum 队列在回 confirm 前已提交到多数派副本,才给出多数工程师以为自己早就有的那种保证。confirm 是对账工具,不是持久化级别本身。
与下一节的关系:confirm 守住了消息进 broker 这一段。消息推给消费者之后,能不能不丢、会不会把消费者撑爆,是另外两个接缝。
2.2消费端不丢、不撑爆:consumer ack 与 prefetch
手动 ack 让 broker 在消费者真正处理完后才删除消息;prefetch 限定一次最多推几条未确认消息,防止单个消费者吞掉整条队列。
两个独立的风险。其一,自动 ack(automatic ack,俗称 ack-on-send)让 broker 在推送瞬间就删消息——消费者处理到一半崩溃,这条已经没了,连 requeue 的机会都没有。手动 ack 把删除推迟到"业务处理完成"那一刻。其二,broker 默认会尽量多推未确认消息给消费者,一个慢消费者会把整条队列的消息全拉进自己内存,撑爆自己(OOM),同时让别的消费者空闲——prefetch(通过 basic.qos 设置)给"未确认消息数"设了上限。
底层机制(比文档深一层):prefetch 限的是一条 channel/consumer 上 Unacked 状态(见 1.4)的消息条数上限。broker 推一条,Unacked 计数 +1;消费者 ack 一条,计数 −1,broker 才补推下一条。所以 prefetch 和 ack 是同一个反馈环的两端:忘了 ack,Unacked 计数只增不减,撞到 prefetch 上限后 broker 停止投递——消费者就此"挂起",但它其实没死,只是再也拿不到新消息。basic.nack 是 RabbitMQ 的扩展,支持 multiple(一次否定到某 tag 为止)和 requeue(退回还是丢弃);basic.reject 功能相同但只能一次一条。注意 4.x 里全局 QoS(global prefetch)已弃用,prefetch 应按 per-consumer 设置。
prefetch 取值有量纲化的经验法则:prefetch ≈ 往返时延 / 单条处理时间。极端值 prefetch=1(入门教程"公平分发"的默认)把吞吐锁死在"每个网络往返才处理一条"——若 RTT 为 125ms,上限约 8 条/秒。实战起步给 10–50,消费者慢、消费者多时调小。
| 设置 | 机制 | 后果 |
|---|---|---|
| 自动 ack | 推送瞬间即删除 | 吞吐最高,消费者崩溃即丢消息,无痕迹 |
| 手动 ack + prefetch 适中 | 处理完才删除,未确认数受限 | 不丢、负载均衡、内存可控——生产默认 |
| 手动 ack + prefetch=1 | 处理完才删除,一次只推一条 | 绝对公平,吞吐被往返时延锁死 |
| 手动 ack + 无限 prefetch | 处理完才删除,但不限未确认数 | 一个消费者吞光队列 → OOM,其余空闲 |
场景走查:四个消费者从一条订单队列取货,单条处理约 50ms,RTT 约 10ms。按经验法则 10/50 远小于 1,prefetch 给个位数即可,取 prefetch=5:每个消费者手里最多压 5 条未确认,broker 在它 ack 后即时补推,四个消费者负载基本均摊。若误设无限 prefetch,启动瞬间最快的那个消费者把上万条全拉进内存,自己 OOM 重启,重启后名下 Unacked 全部 requeue,又一次涌入下一个消费者——形成"击鼓传 OOM"。
channel.basic_qos(prefetch_count=5) # per-consumer,最多 5 条 Unacked
def handle(ch, method, props, body):
try:
process(body) # 业务处理
ch.basic_ack(method.delivery_tag) # 处理完才确认 → broker 删除
except TransientError:
ch.basic_nack(method.delivery_tag, requeue=True) # 可重试:退回队列
except PoisonError:
ch.basic_nack(method.delivery_tag, requeue=False) # 必失败:走 DLX(见 2.4)
channel.basic_consume("orders", handle, auto_ack=False) # 关键:不要自动 ack
症状:消费者偶发"消息凭空消失",日志里没有任何异常。根因:用了 auto_ack=True,broker 在推送的那一刻就删了消息,消费者处理途中崩溃(或被 OOM kill、被部署重启),这条永久丢失、零痕迹。修复:auto_ack=False + 处理成功后再 basic_ack;只有"丢了也无所谓"的指标流才考虑自动 ack。
症状:消费者进程活着、CPU 也不高,却不再消费新消息,队列 Unacked 数顶在某个值不动。根因:业务分支里漏了 basic_ack,Unacked 计数只增不减,撞到 prefetch 上限后 broker 停止投递。看着像"卡死",实则是反馈环被卡住。修复:确保每条消息在所有分支上都恰好 ack/nack 一次。management UI 里"Unacked 持续不降"是这个故障的指纹。
一个消费者设了 prefetch=1,处理每条消息要 100ms,网络 RTT 约 100ms。它的吞吐上限大约是多少?怎么提高?
展开答案(先停 10 秒再点)
约 5 条/秒。prefetch=1 时一条消息的完整周期是"推送(100ms) + 处理(100ms) + ack 回程(100ms)"≈ 200ms 串行往返(推送与 ack 各占半个 RTT),broker 必须等 ack 才推下一条,吞吐被锁在往返节奏上。提高办法:把 prefetch 调到 RTT/处理时间 量级或更高(这里 ≥ 2 就能让处理与网络往返重叠流水线化),让 broker 提前把后续消息推到消费者手里,处理和网络传输并行起来。
与下一节的关系:confirm 和 ack 守的是"在线消息不丢"。但 broker 一旦重启或崩溃,内存里的队列和消息还在不在,是另一套机制——而且这里藏着本章最大的认知陷阱。
2.3重启与崩溃不丢:durable ≠ persistent ≠ safe
durable 是队列定义能否扛重启,persistent 是单条消息要不要落盘,两者都满足也未必扛得住崩溃——那要 quorum。
"重启不丢"被工程师普遍误解为一个开关,实则是三个相互独立的属性,常被混为一谈。durable 是实体属性:队列/exchange 的定义是否在 broker 重启后还存在。persistent(delivery_mode=2)是消息属性:这一条消息是否被标记为要写盘。两者正交——durable 队列里的 transient 消息,重启后队列还在但消息没了;non-durable 队列里的 persistent 消息,重启后队列连同消息一起蒸发。要让一条消息扛过重启,队列 durable 和消息 persistent 必须同时满足,缺一不可。
底层机制(比文档深一层)——本章最重要的一点:满足了 durable + persistent,仍有一个被多数教程跳过的缺口。classic 队列在回送 publisher confirm 之前,并不保证已经 fsync 到磁盘。持久化消息先进的是操作系统的写缓冲,broker 随即回 confirm,真正刷盘是稍后批量进行的(通常有一个 ≤200ms 的窗口)。如果 broker 正好在这个窗口内崩溃(断电、内核 panic、被 OOM kill),这条"durable + persistent + 已 confirm"的消息照样会丢——它从来没真正落到盘上。这就是为什么 quorum 队列存在:quorum 队列基于 Raft,写入要提交到多数派副本之后才回 confirm,单机崩溃由其余副本兜底,给出 durable+persistent 让人误以为自己早已拥有的那种持久性。Raft 的细节留到 03 章 3.4,这里只需记住缺口和补法。
| 属性 | 作用对象 | 挡住的失败 | 挡不住的 |
|---|---|---|---|
| durable | 队列/exchange 定义 | broker 重启后队列定义还在 | 消息本身(若消息非 persistent) |
| persistent | 单条消息 | 消息被标记写盘,配合 durable 扛重启 | fsync 前的崩溃窗口(classic) |
| quorum 队列 | 整条队列 + 多副本 | 单节点崩溃 / 宕机:多数派已落地 | 多数派同时全灭(极端情形) |
场景走查:一个对账系统把队列声明为 durable、消息设 delivery_mode=2、开了 publisher confirms,自认为"零丢失"。压测中给 broker 断电,复盘发现丢了最后约 80ms 的已 confirm 消息。根因不是配置错,而是 classic 队列的 fsync 窗口:那 80ms 的消息回了 confirm 但还在 OS 写缓冲里,断电时一起没了。把队列换成 quorum 队列后重测,同样断电零丢失——因为每条消息在回 confirm 前已写入多数派节点。
症状:队列声明了 durable,重启后队列还在但消息全空。根因:只设了队列 durable,没给消息设 delivery_mode=2,消息是 transient 的,重启时随内存蒸发。反之,给消息设了 persistent 却把队列建成 non-durable,重启后队列定义都没了,消息更无处安放。两个开关是正交的,必须同时打开,且这只防"重启"不防"崩溃窗口"。
一条消息设了 delivery_mode=2(persistent),但发往一个 non-durable 的 classic 队列。broker 正常重启后,这条消息还在吗?
展开答案(先停 10 秒再点)
不在。non-durable 队列的定义本身重启就消失了,承载它的队列都没了,消息标没标 persistent 已无意义——persistent 只决定"消息要不要写盘",不决定"队列要不要保留"。这正是两个属性正交的体现:要扛重启,队列 durable 和消息 persistent 缺一不可;要再扛崩溃窗口,还得上 quorum。
与下一节的关系:消息存住了,但如果某条消息本身有毒——每次处理都失败——nack 重回队列会变成无限循环。把毒消息引开,是下一个接缝。
2.4毒消息不死循环:dead-letter、TTL、重试上限
dead-letter exchange 把被拒绝、过期或超限的消息引到另一个 exchange,而不是丢弃或无限重投;delivery-limit 给重试次数封顶。
手动 ack 解决了"处理失败的消息别删",但若用 nack(requeue=true) 把一条永远会失败的消息退回队列,它会被立刻重新投递、再次失败、再次退回——一个无限重投循环,把 CPU 打满,还卡在队头阻塞后面所有正常消息(head-of-line blocking)。需要一个出口:让"反复失败的"和"过期的"和"积压超量的"消息离开主队列,进到一个专门的地方,等人工或自动重试逻辑处理。这个出口就是 dead-letter exchange(DLX,死信交换机)。
底层机制(比文档深一层):一个队列可以声明 x-dead-letter-exchange 指向一个 DLX。三类事件会让消息被死信化(dead-lettered)而非留在原队列:① 被 nack/reject 且 requeue=false;② TTL 到期(x-message-ttl 设在队列上、或单条消息设 expiration);③ 超过队列长度上限 x-max-length 被挤出。被死信化的消息按原 routing key(或 DLX 配置的覆盖 key)路由进 DLX 绑定的队列。每死信一次,broker 在消息的 x-death 头里累加一条记录(次数、原因、时间、来源队列)——这是排查"这条消息怎么进的死信队列"的取证依据。quorum 队列还有一个内建的 delivery-limit:4.0 起默认值为 20,一条消息被重投 20 次后自动死信化,从机制上掐断毒消息的无限循环,无需手工拼 DLX 计数逻辑。
x-death 头累加计数,配合 delivery-limit 才不会变成换皮的无限循环。场景走查:订单队列消费时,某条消息因下游一个字段 schema 不兼容而每次反序列化都抛异常。消费者对它 nack(requeue=false),消息进 DLX。DLX 绑了两条路径:可重试错误(如下游瞬时 503)进"重试队列",该队列设 30s TTL,过期后再死信回主队列重试,x-death 计数每轮 +1;不可重试错误(如这次的 schema 不兼容)直接进"死信队列"等人工。即便误把它当可重试,quorum 队列的 delivery-limit=20 也会在第 20 次后强制把它打入终态,不会无限打转。
{
"queue": "orders",
"arguments": {
"x-queue-type": "quorum",
"x-dead-letter-exchange": "orders.dlx",
"x-message-ttl": 600000,
"x-max-length": 100000,
"x-delivery-limit": 20
}
}
症状:一条消息在队列里反复出现、CPU 持续打满,后面正常消息迟迟得不到处理。根因:对一条必然失败的消息用了 nack(requeue=true),它立刻回到队列被重投、再失败,形成死循环,还在队头阻塞后续消息。修复:失败时改用 requeue=false 送 DLX,并给 quorum 队列保留默认 delivery-limit(或在 classic 上用 DLX + 重试队列计数)封顶重试次数。
给主队列设了 DLX,但消费者一直用 nack(requeue=true) 重投失败消息。这条毒消息会进 DLX 吗?
展开答案(先停 10 秒再点)
不会——除非命中别的死信条件。死信化只在 requeue=false、TTL 过期、超长或(quorum)超过 delivery-limit 时触发。requeue=true 表示"放回去再试",根本没把消息交给 DLX,它会原地无限重投。DLX 配了不等于自动生效,得让消息真正满足某个死信条件。在 quorum 队列上,即使消费者执意 requeue=true,delivery-limit=20 仍会兜底把它死信化。
与下一节的关系:前四节都在补"消息不丢"的接缝。但可靠性本身有反作用力——队列无限堆积会触发 broker 的自我保护,反而把整个集群的发送端按下暂停键。
2.5可靠性的反作用力:内存/磁盘告警与背压
为了不让自己被撑爆,broker 在内存或磁盘逼近水位线时会阻塞所有生产者——看起来像 publish 卡住,其实是设计中的背压。
把消息可靠地存住,前提是 broker 自己不被存爆。一条没人消费的队列会一直堆,内存与磁盘是有限的。broker 设了告警水位(resource alarm):内存用量逼近高水位(默认约总内存的 60%)或磁盘剩余跌破阈值时,broker 主动阻塞所有发布连接——这是自保的背压,不是故障。消费者照常消费,队列被排空、资源回落到水位线下后,发布自动解除阻塞。可靠性和吞吐在这里短兵相接:与其 OOM 崩掉丢光所有消息,不如先把入口关上。
底层机制(比文档深一层):触发内存告警时,broker 在连接层停止从所有发布连接读取数据——TCP 接收窗口不再推进,生产者的 basic.publish 因此阻塞(或在异步客户端里收到 channel 的 blocked 通知)。这里有一个和 1.1 呼应的细节:阻塞作用在连接层级,所以一条连接上的所有 channel 一起被卡,不是只卡发往满队列的那一个;而且告警是集群级的——任一节点触发,整个集群的发布全停。消费连接不受影响,继续放水。这造成一个极具迷惑性的现象:发送端集体"卡住",消费端却一切正常。
队列实现也影响堆积时的内存行为:经典的 lazy 模式(消息尽量留在磁盘、不常驻内存)和 quorum 队列在大积压下对内存更友好,能延后撞上水位线的时刻;但它们只是推迟而非取消告警——根治还得给队列设上界。
症状:所有生产者的 publish 同时"卡住"或超时,而消费端日志显示一切正常——极易误判成网络或生产者代码问题。根因:某条队列无人消费、无上界,一路堆到 broker 内存高水位(约 60%),触发集群级资源告警,broker 阻塞了所有发布连接。修复:给队列设 x-max-length 或 x-message-ttl 封顶并配 DLX 泄流,用 quorum 或 lazy 行为压低堆积时的内存占用,并监控内存水位告警。把"publish 卡住"的第一反应从"查网络"改成"查队列堆积与资源告警"。
场景走查:一次发布上线把某个消费者服务的镜像配错,消费者全挂,但生产者还在全速 publish。一条队列在十几分钟内堆到几百万条,broker 内存撞上 60% 高水位,触发告警,集群所有发布连接被阻塞。值班同学看到"所有服务 publish 超时",先怀疑网络和 broker 宕机,排查许久才在 management UI 看到内存告警和那条暴涨的队列。修复消费者镜像、队列被消费排空后,内存回落,发布阻塞自动解除。事后给该队列补了 x-max-length 和 DLX——再发生时超量消息进死信而非无限堆积,不再拖垮全集群。
线上所有生产者的 publish 同时变慢甚至卡住,但消费者一切正常、网络也没问题。最先该怀疑什么?
展开答案(先停 10 秒再点)
broker 资源告警触发了背压。"发布端集体阻塞、消费端正常"是 RabbitMQ 内存/磁盘告警的典型指纹——broker 在连接层停止读取发布连接,且这是集群级的,所以全员 publish 卡住。去 management UI 看内存/磁盘水位和有没有暴涨的队列,而不是先查网络。根因通常是某条队列无界堆积;根治是给队列加 x-max-length/TTL 上界并配 DLX 泄流。
本章五节连起来是一条因果链:每加一道"不丢"的保险,都要付出吞吐、延迟或资源的代价。confirm/ack 增加往返,persistent 与 quorum 增加写放大,无界堆积换来的"不丢消息"最终以"全集群发布阻塞"的形式反噬。可靠投递的工程问题从来不是"能不能不丢",而是"为这条消息值得付多少代价、把界划在哪里"。这条权衡线会贯穿 03 章的内部机制和 04 章的选型。
§本章 self-check
先合上教程,把答案写在纸上或编辑器里。写完再点开对照——直接点开等于把这一节当再读一遍。
- 队列声明了 durable、消息设了 persistent、还开了 publisher confirms 并收到了 ack。在 classic 队列上,这条消息还可能丢吗?为什么?换成 quorum 队列结果有何不同?
- publisher confirms 和 consumer ack 都叫"ack",方向和语义各是什么?只开其中一个,会在链路哪一端留下丢失缺口?
- 一个消费者进程活着、CPU 不高,却不再消费新消息,队列 Unacked 数顶住不降。最可能的根因是什么?它和 prefetch 的关系是什么?
- (设计题)给一条核心订单队列设计一套"消息不丢 + 毒消息不死循环 + 不拖垮集群"的完整方案,列出你会设置的队列类型与关键参数,并说明每一项挡的是哪种失败。
答案(先做完再展开)
- classic 队列上仍可能丢:它在回 confirm 前并不 fsync,消息还在 OS 写缓冲里,broker 在这 ≤200ms 窗口内崩溃就丢。confirm 只代表"broker 按它的承诺处理完",对 classic 这个承诺不含落盘。quorum 队列在回 confirm 前已提交到多数派副本,单机崩溃由其余副本兜底,因此不丢。
- publisher confirms 是 broker → producer 的回执,语义是"消息已被 broker 托管";consumer ack 是 consumer → broker 的回执,语义是"已处理完,可删除"。只开 confirm 不开手动 ack,消费端崩溃丢消息;只开手动 ack 不开 confirm,发送端进 broker 失败时无感知。两段方向相反、互不替代。
- 根因是漏了 ack。Unacked 计数只随投递增加、不随 ack 减少,撞到 prefetch 上限后 broker 停止投递,消费者看着像假死,实则反馈环被卡住。prefetch 限的就是 Unacked 上限,所以"忘 ack"必然以"Unacked 顶在 prefetch 值不动"的形式暴露。
- 参考方案:队列用 quorum(多副本,挡单机崩溃,且自带
delivery-limit=20挡无限重投);消息设delivery_mode=2persistent + 队列 durable(挡重启);生产端开 publisher confirms +mandatory(挡发送端丢失与无路由静默丢弃);消费端手动 ack + 适中prefetch(挡消费端崩溃丢失与撑爆);配 DLX 接收requeue=false/TTL/超长的消息(挡毒消息与积压);给队列设x-max-length/x-message-ttl上界(挡无界堆积触发全集群发布阻塞)。每一项对应一种具体失败,缺哪项就在哪处留缺口。
用 TTL + DLX 拼一个"延迟重试 30 秒"
RabbitMQ 没有原生的延迟队列。要求:消费失败的消息不立即重投,而是等 30 秒再回到主队列重新消费,且最多重试 5 次,超过就进终态死信队列。手上只有普通队列、TTL、DLX 三样积木——画出队列与 exchange 的拓扑,说明消息怎么"绕一圈"实现延迟,以及 5 次上限靠什么字段判定。
提示(卡住再展开)
关键技巧:让消息在一个没有消费者的"等待队列"里靠 TTL 自然过期,过期即死信。拓扑:主队列消费失败 → nack(requeue=false) 进 DLX-A → 路由到"等待队列"(设 x-message-ttl=30000,且它的 DLX 指回主队列对应的 exchange)→ 30s 后 TTL 到期,消息被死信回主队列重试。等待队列没有消费者,消息只能靠过期离开,"延迟"就来自这段 TTL。重试次数读消息 x-death 头里的累计计数,超过 5 就在消费端改判进终态死信队列(quorum 队列也可直接靠 delivery-limit 兜底,但那是次数上限不是延迟)。想一想:如果不同消息要不同延迟(30s/5min/1h),单条 TTL 的"队头阻塞"会带来什么问题?