RabbitMQ
一、是什么
1. RabbitMQ 转移的是消息责任,不是业务事务
RabbitMQ 是一个消息代理。生产者不直接寻找每个消费者,而是把消息发布给 exchange;exchange 按 binding 选择一个或多个 queue;queue 保存尚未完成的 delivery;消费者处理完成后再确认。它把“谁现在负责保管和推进这条消息”从应用进程中分离出来,却不会自动把订单数据库、库存数据库和外部接口变成一个原子事务。
订单服务记录“发送完成”,库存却没有变化,可能发生在不同位置。消息可能尚未离开生产者进程,连接已经断开;broker 也可能接受了 publish,但 routing key 没匹配任何 queue。即使消息已经进入 queue,也可能没有在线消费者;delivery 已交给消费者时会处于 unacked,而业务事务仍未提交;数据库已经提交后,消费者若在 ack 到达 broker 前退出,消息还会再次投递。
这些位置对应三次相互独立的责任转移。Publisher confirm 说明 broker 是否接管了 publish;exchange、binding 和 queue 决定消息是否拥有实际去处;consumer acknowledgement 说明当前 delivery 是否可以从 queue 的待处理状态中结束。
Confirm 不表示消费者处理成功,consumer ack 也不证明数据库一定正确。网络在 confirm 返回前中断时,生产者甚至无法判断 broker 是“没有收到”还是“已经接管但响应丢失”。RabbitMQ 的可靠性指南因此要求发布者处理未确认消息,消费者处理重投,队列类型再承担相应的数据安全责任。
2. Connection、Channel 与 vhost 组成运行边界
2.1 Connection 承担网络生命周期
AMQP 0-9-1 connection 是一条长期 TCP 连接,负责认证、心跳、流控通知和连接恢复。每次发布都新建 connection 会放大 TLS、认证与 socket 开销;长期连接若没有 heartbeat 和断线恢复,又会把半开连接拖到业务超时才暴露。
连接成功只表示客户端进入了某个 vhost。它不代表 exchange 或 queue 存在,也不代表当前身份拥有发布、消费或声明资源的权限。
2.2 Channel 是协议会话和确认作用域
多个 channel 可以复用一条 connection,分别承担发布、消费和拓扑声明。Publisher confirm 的序号、consumer delivery tag 和事务状态都属于 channel。收到 delivery 的 channel 才能 ack 对应 delivery tag;把 channel 跨线程并发使用,或者重连后继续使用旧 tag,通常会得到 unknown delivery tag 并关闭 channel。
连接恢复之后,channel、confirm 模式和 consumer subscription 都需要重新建立。自动重连只修复传输层,不能假定上一条未收到 confirm 的 publish 没有进入 queue。
2.3 vhost 与权限隔离资源
vhost 是 exchange、queue、binding、policy 和权限的命名空间。相同名称的 orders.events 可以存在于不同 vhost,彼此没有路由关系。RabbitMQ 资源权限由 configure、write、read 三个正则组成。configure 控制声明、删除和改变资源,write 控制向 exchange 发布以及某些绑定操作,read 控制从 queue 消费、获取或清理消息。
Management UI 的用户 tag 决定管理界面能力,不替代 vhost 资源权限。应用账号不应因为需要消费就同时获得拓扑删除和 policy 变更能力。
3. Exchange、Binding 与 Queue 共同决定路由
生产者向 exchange 发布消息,并提供 routing key。Exchange 本身通常不保存普通 AMQP 消息,而是根据类型和 binding 计算目标 queue:
| Exchange 类型 | 匹配方式 | 典型用途 | 常见误用 |
|---|---|---|---|
direct | routing key 精确匹配 | 明确命令、固定任务类型 | 为每个实例创建队列,意外变成多份副本 |
topic | * 匹配一个词,# 匹配多个词 | 按领域、地区、版本订阅事件 | routing key 没有稳定词汇表,规则难以推断 |
fanout | 忽略 routing key,复制到所有绑定队列 | 缓存失效、广播通知 | 多个消费者共用一个队列,实际变成竞争消费 |
headers | 根据消息头匹配 | 少量无法用 routing key 表达的规则 | 把复杂业务判断搬进 broker |
同一 queue 上的多个消费者是竞争消费,一条 delivery 在同一时刻只交给其中一个消费者。多个独立 queue 绑定同一个 exchange,则每个 queue 各自保存一份消息。订单投影、通知和风控都需要收到 order.created 时,应建立三个订阅队列;让三类服务共用一个队列,只会让它们互相抢消息。
空名称的 default exchange 会把 routing key 等于 queue 名的消息直接路由到该 queue,适合最小实验。正式业务更适合显式 exchange 和 binding,因为路由关系需要独立演进、授权与诊断。
当 publish 没有匹配任何 queue 时,mandatory=false 允许 broker 丢弃它。启用 mandatory 后,broker 会用 Basic.Return 把不可路由消息退给生产者。Publisher confirm 与 return 仍是两件事:一个不可路由 publish 也可能获得 confirm,因为 broker 已经完成了对它的处理;生产者必须同时处理 return 和 confirm。
4. 一条消息在 Queue 中经历什么
消息路由到 queue 后先进入 ready。Broker 选择有投递空间的 consumer,把消息交给它后转为 unacked。Consumer 执行 ack 表示处理完成,结束当前 delivery;执行 reject 或 nack(requeue=false) 表示不再回到原队列,消息按 DLX 配置转移或丢弃;执行 nack(requeue=true) 会返回原队列,可能很快再次投递。Channel 或 connection 断开时,所有未确认 delivery 也会自动返回队列。
因此 ready 很高通常表示没有可用消费者或净消费能力不足;unacked 很高则更接近消费者处理慢、prefetch 过大、线程阻塞或忘记 ack。只看总消息数无法区分这两类问题。
prefetch_count 限制每个消费者允许保留多少条未确认 delivery。它只影响订阅式消费,不影响 basic.get 这种轮询式拉取。Prefetch 太小会增加等待往返,太大会把大量消息锁在慢消费者内存中,并在进程退出时形成重投尖峰。
5. Classic Queue、Quorum Queue 与 Stream 保存不同状态
5.1 Classic Queue
Classic queue 适合可重建、短生命周期、低复制要求或高队列周转率的消息。单节点 classic queue 的内容不因 RabbitMQ 集群而自动复制;旧式 classic mirrored queue 已从 RabbitMQ 4.x 移除。需要复制的数据不能再沿用旧 HA policy 思路。
5.2 Quorum Queue
Quorum queue 使用 Raft 复制入队、delivery 和 acknowledgement 等队列状态。Leader 把操作写入日志并复制到多数成员后,才满足安全确认条件。三成员可以容忍一个成员不可用;同时失去两个成员就无法形成多数派,队列停止推进以避免两个分区各自接受不一致写入。
Quorum queue 适合订单、支付指令等需要复制和明确故障边界的长生命周期队列。它必须 durable,不适合 exclusive 临时队列、高频创建删除、只追求最低延迟或数百万条以上的长期积压。官方 Quorum Queue 指南建议在这些场景考虑 classic queue 或 stream。
5.3 RabbitMQ Stream
Stream 是可复制、追加式、按保留策略清理的日志,消费者可以从不同 offset 重读。它更适合大扇出、较长保留和重复读取,不等同于普通 queue 的 ready/unacked/逐条删除模型。若需求已经变成长期事件日志和多次回放,应直接评估 Stream 或 Kafka,而不是无限放大 quorum queue。
队列类型在声明时确定,不能靠 policy 把已有 classic queue 原地变成 quorum queue。迁移需要创建新队列、切换路由或双写,并明确旧消息怎样排空。
6. 重试、Delivery Limit 与 Dead Letter 是不同机制
瞬时错误和永久错误不能走同一条路径。单行锁冲突、单租户限流可以延迟后重试;消息格式错误、缺少必填字段通常不会因为等待而恢复;数据库整体离线时,让每条消息独立重试只会制造重投风暴,此时应暂停消费者并保护积压。
Quorum queue 会记录失败投递次数,并可通过 delivery limit 把毒消息送往 dead letter exchange。RabbitMQ 4.3 还支持 quorum queue 内部的线性延迟重试,消息不必在一组 TTL queue 之间来回改写。延迟由 delayed-retry-min、delayed-retry-max 和 delivery count 计算,适合局部、可恢复的失败。RabbitMQ 4.3 说明同时区分了 returned 与 failed:AMQP 0-9-1 的 basic.reject(requeue=true) 会把本次尝试标为失败、增加 delivery count,并可进入 failed 延迟;basic.nack(requeue=true) 属于普通 returned,不增加 delivery count,也不会受 delivery limit 约束。应用必须按实际协议动作选择策略,不能把 reject 和 nack 当成同义方法。
DLX 是另一个 exchange。消息被拒绝、过期、超过长度或 delivery limit 时,源 queue 可以把它重新发布到 DLX。默认 dead lettering 可能在目标不可用时丢失;关键 quorum queue 若需要 at-least-once dead lettering,要显式配置对应策略、overflow=reject-publish,并确保目标 exchange 和 queue 可用。即便如此,网络和确认窗口仍可能产生重复,死信消费者一样需要幂等。
7. 官方资源、版本、Erlang 与协议边界
RabbitMQ Server 是运行在 Erlang/OTP 虚拟机上的消息代理,不是一个脱离运行时的单文件程序。安装 Server 前必须先确认 Erlang 兼容矩阵;集群中还应让全部节点保持相同的 Erlang 大版本,并尽量使用完全相同的补丁版本。常用的一手入口如下:
| 资源 | 地址 | 主要用途 |
|---|---|---|
| 官网与 4.3 文档 | rabbitmq.com / RabbitMQ 4.3 Documentation | 产品概念、配置和运维 |
| 安装入口 | Installing RabbitMQ | Linux、Windows、macOS、容器和 Kubernetes 安装路径 |
| 版本与支持周期 | Release Information | 当前补丁、支持状态和 release notes |
| Erlang 兼容矩阵 | Erlang Version Requirements | RabbitMQ 与 Erlang/OTP 的最小、最大兼容版本 |
| 源码与发布制品 | rabbitmq-server | 源码、标签、安装包和签名 |
| 客户端列表 | Client Libraries | 按语言和协议选择官方或社区客户端 |
本文固定 RabbitMQ 4.3.5、Erlang/OTP 27.x 和 Pika 1.3.x 作为可复现实验基线。RabbitMQ 4.3.5 要求 Erlang/OTP 至少为 27.0,完整支持 27.x;Erlang 28 只对新集群部分支持,Erlang 29 不受 4.3 支持。不能因为操作系统仓库里存在 rabbitmq-server 就直接安装,旧发行版仓库常同时带来过期的 RabbitMQ 与 Erlang。
镜像固定为 rabbitmq:4.3.5-management,不要把 latest 或浮动 minor 标签当作可回滚制品。生产升级还要分别核对 Server、Erlang、客户端、插件、feature flags 和元数据迁移路径;“容器能启动”不代表这些组合受支持。
RabbitMQ 4.3 只使用基于 Raft 的 Khepri 保存 vhost、用户、权限、exchange、queue 定义、binding、policy 和 runtime parameter 等元数据;普通 queue 或 stream 中的消息由各自的数据结构保存。Definitions 导出的是前一类元数据,不会把后一类消息一起导出。一个节点的主要运行关系可以概括为:
Erlang VM 负责大量轻量进程、调度、网络和分布式通信;RabbitMQ application 在其上实现连接、路由、队列、复制与插件。节点名标识集群成员,Erlang cookie 用于节点之间和本机 CLI 的认证,它不是应用登录密码。插件会增加协议、管理或监控入口,但未启用的插件不会因为安装包存在就自动监听端口。
RabbitMQ 支持 AMQP 0-9-1、AMQP 1.0、MQTT、STOMP 和 Stream 等入口。本文的 exchange、publisher confirm、delivery tag 与 manual ack 使用 AMQP 0-9-1。未加密 AMQP 常用 5672,TLS AMQP 常用 5671,Management HTTP 默认 15672,Prometheus 插件默认 15692,集群分布式通信还会使用 25672 和 EPMD 4369。端口不是都应对外开放:业务网络只开放所需协议端口,管理、监控和节点间端口按来源网络隔离。管理页面能打开,也不代表应用连接的 AMQP 端口、TLS、vhost 和资源权限正确。
8. Node、Cluster、声明参数、Policy 与 Feature Flag
一个 RabbitMQ node 是一个带唯一 node name 和数据目录的 Erlang 节点;多个 node 组成 cluster 后共享元数据,但消息是否复制仍由 queue 类型和成员数决定。三节点集群中的单副本 classic queue 仍可能随它所在节点永久丢失,不能把“节点已入集群”当成“所有消息已有三份副本”。Broker 通常指对客户端提供协议入口和消息能力的整个服务角色,在单节点环境里常与 node 混用;排障时必须落到具体 node、具体 queue leader 和具体副本成员。
Queue 声明会同时确定名称、durable、exclusive、auto-delete、类型和 arguments。durable 表示重启后保留定义,不等于每条消息都持久;消息还需要 persistent delivery mode,并由对应 queue 的写入与复制语义接管。exclusive queue 绑定到声明它的 connection,连接关闭时删除;auto-delete 需要至少有过 consumer,最后一个 consumer 消失后才触发删除。服务端命名的临时 queue 适合请求回复或实例级通知,业务共享队列应使用稳定、版本化名称。
Policy 把 TTL、长度、DLX、delivery limit 等可变规则按名称正则施加到一组资源;operator policy 由运维给出不可被应用放大的保护上限;声明参数属于对象本身。同一个数值键在多层同时出现时,RabbitMQ 通常采用更保守的结果,但不可变属性冲突会关闭 channel 并返回 PRECONDITION_FAILED,不能靠更高 policy 优先级覆盖。
Feature flag 控制集群内需要版本协调的能力。升级前要让旧版本的 stable flags 全部启用,升级后只有在所有节点稳定、确认不再回到旧版本时才启用新 flags。Feature flag 一旦启用通常不能关闭,因此它既是能力开关,也是升级回退边界。
二、为什么
1. 先判断传递的是任务、命令、事件还是日志
消息中间件选型从业务责任开始,而不是从吞吐宣传开始:
| 消息性质 | 核心要求 | RabbitMQ 中的自然模型 |
|---|---|---|
| 后台任务 | 一份任务由一个 worker 完成,可失败重试 | 一个 queue + 竞争消费者 + manual ack |
| 业务命令 | 指定能力处理,重复必须收敛 | direct exchange + quorum queue + 幂等键 |
| 领域事件 | 多个下游各自处理,生命周期相互独立 | topic exchange + 每个下游独立 queue |
| 在线通知 | 低延迟优先,离线可丢 | fanout 或临时 classic queue |
| 可回放日志 | 长期保留,多组消费者从不同位置读取 | RabbitMQ Stream、Kafka 或 Pulsar |
RabbitMQ 的优势是路由、逐条确认、竞争消费、有限积压和失败隔离可以直接映射到业务处理。若主要需求是按 offset 保留数周、批量扫描历史或让大量消费组反复回放,普通 queue 的删除语义反而会增加复杂度。
2. 从接收者关系推导 Exchange 与 Queue
“一个事件有几个消费者”必须拆成两问:有几个业务订阅者,以及每个订阅者需要多少处理实例。
订单创建后,库存、通知、风控各自都要处理,应建立三个 queue,并分别绑定 order.created。库存服务部署四个实例时,这四个实例共同消费库存 queue;它们提高并发,不会让每个实例各收一份。只有广播到每个实例的本地缓存失效通知,才需要每个实例拥有自己的临时 queue。
Routing key 应表达稳定领域分类,而不是服务器名或消费者版本。Topic exchange 的层次可以采用 领域.实体.动作.版本,例如 commerce.order.created.v1。版本变化若不兼容,应显式新建 routing key 或 queue,不能让消费者靠猜测 payload 字段适配。
3. 从数据角色选择 Queue 类型
选择 classic、quorum 或 stream 时,要同时回答消息能否从事实库重建、单节点永久损坏时允许丢多少、最大积压会持续多久、消息只处理一次还是需要多次回放,以及系统能否承担副本带来的磁盘、网络和确认延迟。
图片缩略图任务可以从对象存储清单重建,短暂单机队列可能足够。订单扣库存命令若无法无损重建,更适合三成员 quorum queue。审计事件需要多组消费者长期回放时,应使用 stream 或提交日志。把全部 queue 默认设为 quorum 会浪费资源,把全部 queue 默认设为 classic 又会把关键数据暴露在单节点故障中。
4. 从不确定窗口设计可靠交付
RabbitMQ 常见的“至少一次”来自两个窗口。Publish 已被 broker 接管而 confirm 在网络中丢失时,生产者会重发;数据库已经提交而 ack 尚未到达 broker 时,消费者退出会触发 broker 重投。
生产者需要稳定 event_id 和 transactional outbox。订单与 outbox event 在同一个数据库事务中提交;relay 只有在 publish 获得 confirm 且没有 return 后才把 outbox 标记为已发布。超时保留 pending,并用同一 event ID 重发。
消费者需要 inbox 或业务唯一约束。它在一个数据库事务中记录 event_id 并完成库存预留,提交后再 ack。若同一消息重投,唯一约束让第二次处理成为无副作用的重复,而不是再次扣减。
Broker 参数不能消除外部数据库双写。XA 可以覆盖少数受控系统,却把协调器、恢复和长事务带进链路;大多数事件驱动系统更适合 outbox、幂等和补偿。
5. 从失败性质设计 Retry 与 Dead Letter
重试前先给错误分类。连接闪断、行锁冲突和单租户 429 适合有限次数、带退避的重试;数据库整体离线时应暂停消费,恢复依赖后继续;schema 不合法或业务对象永久不存在时应直接死信;下游已经执行但返回未知时,应先按业务幂等键查询结果,再决定是否重试。
requeue=true 没有天然延迟,可能让同一条消息在消费者之间高速循环。Quorum delayed retry 适合 RabbitMQ 4.3 的短中期退避;需要日历时间、人工审批或跨天编排时,应使用调度器或工作流系统。死信不是垃圾桶:它需要明确失败原因、修复方式、限速重放和最终清理周期。
6. 从故障责任选择部署形态
| 形态 | 适用场景 | 仍由应用承担的责任 | 主要限制 |
|---|---|---|---|
| 本地单容器 | 协议学习、契约测试 | confirm、ack、幂等、清理 | 无法验证副本与多数派 |
| 共享开发实例 | 多项目联调 | vhost、权限、命名、配额 | 噪声和误操作会相互影响 |
| 三节点自建集群 | 关键内部消息、自主管理 | 客户端恢复、业务一致性 | Erlang/RabbitMQ 升级、磁盘、网络和值守 |
| 托管多可用区 | 希望转移节点生命周期 | 消息模型、权限、重试、降级 | 配额、能力差异、网络和供应商故障 |
RabbitMQ 集群适合可靠 LAN。跨 WAN 的延迟和分区会直接放大 Raft 与 confirm;跨地域连接通常使用独立集群加 Federation、Shovel 或应用级复制,而不是把一个集群横跨两个机房。
7. 顺序、容量和成本必须一起决定
一个 queue 加一个 active consumer 最容易保持投递顺序,却限制吞吐和可用性。增加消费者会并发处理,完成顺序不再稳定。Single Active Consumer 可以让 quorum queue 在 active consumer 失效后切换,但重投仍可能让业务看到旧事件;真正按订单有序时,还要用业务键分片和状态版本拒绝过期更新。
积压恢复取决于净消费能力:
净排空速率 = 总消费速率 - 到达速率
预计恢复时间 = backlog / 净排空速率生产每秒 800 条、消费每秒 900 条时,十万条积压至少需要一千秒;若消费不高于到达,扩磁盘只能推迟故障。容量还要包含消息大小分布、重试放大、prefetch、quorum 副本数、磁盘高水位和恢复期间的业务峰值。
8. 什么时候不应继续使用 RabbitMQ
如果数据要保留数天到数月,并让多个消费组从不同 offset 回放,或者单条记录需要被海量订阅者重复读取,问题已经超出普通 queue。业务依赖跨 partition 的流处理、窗口聚合或大规模日志生态,消息积压长期达到数百万并持续增长,以及系统需要跨地域自治而不是低延迟局域网集群,也都是更换模型的信号。另一种相反的误用是实际只需要同步查询,却用消息绕过明确的 API 失败语义。
这时应分别评估 RabbitMQ Stream、Kafka、Pulsar、NATS 或同步调用。选择不同工具不是否定 RabbitMQ,而是让消息生命周期与产品模型一致。
9. 用四个场景校准选择
图片转码任务只需要一份工作由一个 worker 完成,允许失败重试,源图片又能重新扫描,适合一个 classic 或 quorum queue 加多个竞争消费者。是否选 quorum 取决于任务重建成本和单节点故障时允许丢失多少,而不是“生产一律 quorum”。
订单创建后,库存、通知和风控都必须得到事件,应该让三个独立 queue 绑定同一个 topic exchange;每个服务再用自己的多个实例竞争本服务 queue。把三类服务都连到一个 queue 会让它们互相抢走消息,不是广播。
同一账户的状态更新要求有序时,单 queue 多消费者也不能保证完成顺序。可以按账户分片到稳定 queue、使用 Single Active Consumer 限制活动消费者,并让业务状态携带单调版本拒绝旧事件。RabbitMQ 只能保证一个 queue 中的投递约束,不能替业务建立跨 queue、重投和外部事务的全局顺序。
审计日志需要保存较长时间,让多个团队从不同位置反复回放时,普通 queue 的 ack 后删除模型并不自然。此时先评估 RabbitMQ Stream;若还需要大规模分区生态、流处理和更长保留,再评估 Kafka 或 Pulsar。不要为了保留历史而让 quorum queue 长期堆积到无法恢复。
无论场景名称是什么,最终都要回答:消息能否重建,允许丢多少,是否需要广播或回放,顺序按什么键成立,最长积压多久,失败由谁重试,RPO/RTO 是多少,消费副作用怎样幂等,团队能否维护 TLS、复制、监控和升级。缺少这些答案时,exchange 和 queue 类型只是猜测。
三、怎么做
1. 先选版本,再完成第一次安装
1.1 在安装前确认运行方式和兼容组合
第一次学习优先使用 Docker Compose,它能固定 RabbitMQ、插件和数据卷,清理边界也最明确。需要理解真实服务目录、systemd、日志和主机限制时,再走 Ubuntu/Debian 原生安装。生产环境通常采用三节点 Linux、Kubernetes Operator 或托管服务,不能把单容器当作高可用方案。
| 目标 | 推荐入口 | 此时先不解决什么 |
|---|---|---|
| 十分钟看到第一条消息 | Docker Compose 单节点 | 节点故障、TLS、容量和恢复 |
| 学习真实服务管理 | Ubuntu/Debian 官方 apt 仓库 | 多数派与跨节点复制 |
| 验证 Quorum Queue | 三个独立节点或容器 | 跨地域灾备 |
| 生产上线 | 三个故障域内节点、Operator 或托管多可用区 | 应用幂等、Outbox 和业务补偿仍由项目承担 |
先检查 CPU 架构、系统版本、磁盘和端口。Team RabbitMQ 的 Erlang apt 仓库示例只为 amd64 提供包;arm64 应按官方安装页选择受支持的 Erlang 来源,不能照抄 amd64 源。
uname -m
. /etc/os-release && printf '%s %s\n' "$ID" "$VERSION_ID"
df -h
ss -lnt '( sport = :5672 or sport = :15672 )'本文安装结果应同时满足 RabbitMQ 4.3.5、Erlang/OTP 27.x。版本号相同也不表示数据目录可以任意互换;节点名、cookie、feature flags、插件和元数据状态同样属于运行条件。
1.2 用 Docker Compose 获得可重复的本地节点
创建 .env:
RABBITMQ_IMAGE=rabbitmq:4.3.5-management
RABBITMQ_AMQP_PORT=5672
RABBITMQ_UI_PORT=15672
RABBITMQ_DEFAULT_USER=local_admin
RABBITMQ_DEFAULT_PASS=replace-with-a-random-local-password
RABBITMQ_DEFAULT_VHOST=local_lab创建 compose.yaml:
services:
rabbitmq:
image: ${RABBITMQ_IMAGE:?set RABBITMQ_IMAGE}
container_name: lab-rabbitmq
hostname: lab-rabbitmq
restart: unless-stopped
ports:
- "127.0.0.1:${RABBITMQ_AMQP_PORT:-5672}:5672"
- "127.0.0.1:${RABBITMQ_UI_PORT:-15672}:15672"
environment:
RABBITMQ_DEFAULT_USER: ${RABBITMQ_DEFAULT_USER:?set user}
RABBITMQ_DEFAULT_PASS: ${RABBITMQ_DEFAULT_PASS:?set password}
RABBITMQ_DEFAULT_VHOST: ${RABBITMQ_DEFAULT_VHOST:-local_lab}
volumes:
- rabbitmq-data:/var/lib/rabbitmq
healthcheck:
test: ["CMD", "rabbitmq-diagnostics", "-q", "ping"]
interval: 10s
timeout: 5s
retries: 12
start_period: 20s
volumes:
rabbitmq-data:只绑定 127.0.0.1,避免开发口令和 Management UI 暴露到局域网。启动并分别检查进程、节点和监听端口:
docker compose up -d rabbitmq
docker compose ps rabbitmq
docker exec lab-rabbitmq rabbitmq-diagnostics -q ping
docker exec lab-rabbitmq rabbitmq-diagnostics listeners
docker exec lab-rabbitmq rabbitmqctl await_startup
docker exec lab-rabbitmq rabbitmqctl versionManagement UI 位于 http://127.0.0.1:15672。ping 成功只说明 Erlang 节点能响应 CLI;应用仍可能因为 AMQP 端口、vhost 或权限错误而失败。
环境变量只在空数据目录首次初始化默认用户。已有 volume 保存着旧用户时,修改 .env 不会重置密码。应显式修改用户,不能把持久化状态误判成 Compose 缓存。
1.3 在 Ubuntu 上安装原生服务
下面以 Ubuntu 24.04(发行版本标识 Noble)和 amd64 为例。RabbitMQ 官方明确不建议直接使用 Ubuntu/Debian 自带的旧版本包;这里添加 Team RabbitMQ 仓库,同时把 Erlang 固定在受完整支持的 27.x,把 Server 固定到 4.3.5。其他发行版必须替换为官方安装页对应的仓库代号。
先安装仓库工具,下载并核对 Team RabbitMQ 公钥指纹:
sudo apt-get update
sudo apt-get install -y curl gnupg apt-transport-https
curl -1sLf \
'https://keys.openpgp.org/vks/v1/by-fingerprint/0A9AF2115F4687BD29803A206B73A36E6026DFCA' \
| sudo gpg --dearmor -o /usr/share/keyrings/com.rabbitmq.team.gpg
sudo gpg --show-keys --with-colons /usr/share/keyrings/com.rabbitmq.team.gpg \
| grep 'fpr:::::::::0A9AF2115F4687BD29803A206B73A36E6026DFCA:'指纹必须完整匹配后再写入仓库。两个镜像源提供相同软件包,一个不可达时 apt 可以使用另一个:
sudo tee /etc/apt/sources.list.d/rabbitmq.list >/dev/null <<'EOF'
deb [arch=amd64 signed-by=/usr/share/keyrings/com.rabbitmq.team.gpg] https://deb1.rabbitmq.com/rabbitmq-erlang/ubuntu/noble noble main
deb [arch=amd64 signed-by=/usr/share/keyrings/com.rabbitmq.team.gpg] https://deb2.rabbitmq.com/rabbitmq-erlang/ubuntu/noble noble main
deb [arch=amd64 signed-by=/usr/share/keyrings/com.rabbitmq.team.gpg] https://deb1.rabbitmq.com/rabbitmq-server/ubuntu/noble noble main
deb [arch=amd64 signed-by=/usr/share/keyrings/com.rabbitmq.team.gpg] https://deb2.rabbitmq.com/rabbitmq-server/ubuntu/noble noble main
EOF
sudo tee /etc/apt/preferences.d/rabbitmq-erlang >/dev/null <<'EOF'
Package: erlang*
Pin: version 1:27.*
Pin-Priority: 999
EOF
sudo tee /etc/apt/preferences.d/rabbitmq-server >/dev/null <<'EOF'
Package: rabbitmq-server
Pin: version 4.3.5-1
Pin-Priority: 999
EOF
sudo apt-get update
apt-cache policy erlang-base rabbitmq-serverCandidate 必须落在预期仓库、Erlang 27.x 和 RabbitMQ 4.3.5。确认后安装 Erlang 运行组件与 Server:
sudo apt-get install -y \
erlang-base erlang-asn1 erlang-crypto erlang-eldap erlang-ftp \
erlang-inets erlang-mnesia erlang-os-mon erlang-parsetools \
erlang-public-key erlang-runtime-tools erlang-snmp erlang-ssl \
erlang-syntax-tools erlang-tftp erlang-tools erlang-xmerl \
rabbitmq-server=4.3.5-1
erl -noshell -eval \
'io:format("OTP ~s~n", [erlang:system_info(otp_release)]), halt().'
sudo rabbitmqctl version预期分别看到 OTP 27 和 4.3.5。如果 apt 计划安装 Erlang 28、不同 RabbitMQ 补丁或大量降级包,应先停止并检查 apt-cache policy、仓库代号和 pin,而不是使用 --allow-downgrades 强行继续。
1.4 识别服务、配置、数据、日志和插件
Linux 包以非特权用户 rabbitmq 运行 rabbitmq-server.service。包通常不会替读者生成 rabbitmq.conf;没有该文件不表示安装损坏。常用位置是:
| 对象 | Ubuntu/Debian 默认位置 | 作用 |
|---|---|---|
| 主配置 | /etc/rabbitmq/rabbitmq.conf | listener、TLS、内存/磁盘水位和插件参数 |
| 高级配置 | /etc/rabbitmq/advanced.config | 少数无法用 sysctl 格式表达的 Erlang 配置 |
| 环境配置 | /etc/rabbitmq/rabbitmq-env.conf | node name、目录等启动环境 |
| 节点数据 | /var/lib/rabbitmq/mnesia/ | Khepri、消息存储和节点状态;名称保留历史兼容含义 |
| 节点 cookie | /var/lib/rabbitmq/.erlang.cookie | 节点与 CLI 认证,不是应用密码 |
| 文件日志 | /var/log/rabbitmq/ | RabbitMQ 节点日志;同时结合 journal 查看服务启动 |
启动服务并从进程、节点、端口、日志四层检查:
sudo systemctl enable --now rabbitmq-server
sudo systemctl status rabbitmq-server --no-pager
sudo rabbitmq-diagnostics -q ping
sudo rabbitmqctl await_startup
sudo rabbitmqctl version
sudo rabbitmq-diagnostics listeners
sudo rabbitmq-diagnostics status
sudo journalctl -u rabbitmq-server -n 100 --no-pager
sudo ls -ld /etc/rabbitmq /var/lib/rabbitmq /var/log/rabbitmqactive 只说明 systemd 进程仍在,ping 只说明 CLI 能通过 cookie 访问 Erlang 节点,listeners 才说明 AMQP 端口已绑定;应用认证和消息收发还要继续验证。若 CLI 报 cookie 或节点名错误,不要复制别的集群 cookie 覆盖数据节点,应先查看日志中的 node name、有效用户和 cookie 路径。
Management 是可选插件。开发节点可启用后创建独立管理员并删除默认 guest;生产管理口应留在管理网络或反向代理之后,不能直接暴露公网:
sudo rabbitmq-plugins enable rabbitmq_management
sudo rabbitmq-plugins list --enabled
sudo rabbitmqctl add_user local_admin
sudo rabbitmqctl set_user_tags local_admin administrator
sudo rabbitmqctl set_permissions -p / local_admin '.*' '.*' '.*'
sudo rabbitmqctl delete_user guest
sudo rabbitmq-diagnostics listenersadd_user 不在命令行携带密码时会交互输入。确认 15672 已监听后再访问 Management UI;网页可登录只证明 HTTP 管理入口与管理 tag 可用,不证明 AMQP vhost 权限正确。服务端生产主线以 Linux 为准;其他桌面系统更适合作为客户端或通过容器完成本地实验,避免把桌面服务机制混入服务器运维路径。
2. 发送并消费第一条消息
先完成一个没有重试、集群和数据库的最小闭环。创建隔离环境并安装 Pika:
python3 -m venv .venv-rabbitmq
source .venv-rabbitmq/bin/activate
python -m pip install --upgrade pip
python -m pip install "pika==1.3.2"创建 hello_rabbit.py。脚本使用显式 direct exchange、durable queue 和 binding;publish 负责声明并发送,consume 只取一条并手动 ack。先用 basic_get 是为了看清第一条消息的状态,后面的并发消费会改用 basic_consume。
import os
import sys
import pika
if len(sys.argv) != 2 or sys.argv[1] not in {"publish", "consume"}:
raise SystemExit("usage: python hello_rabbit.py publish|consume")
params = pika.URLParameters(os.environ["AMQP_URL"])
params.heartbeat = 30
params.blocked_connection_timeout = 10
with pika.BlockingConnection(params) as connection:
channel = connection.channel()
channel.exchange_declare(
exchange="learn.events", exchange_type="direct", durable=True
)
channel.queue_declare(queue="learn.hello", durable=True)
channel.queue_bind(
exchange="learn.events", queue="learn.hello", routing_key="hello"
)
if sys.argv[1] == "publish":
channel.basic_publish(
exchange="learn.events",
routing_key="hello",
body=b"hello RabbitMQ",
properties=pika.BasicProperties(
delivery_mode=pika.DeliveryMode.Persistent,
message_id="evt-learn-hello-v1",
content_type="text/plain",
),
)
print("published message_id=evt-learn-hello-v1")
else:
method, properties, body = channel.basic_get(
queue="learn.hello", auto_ack=False
)
if method is None:
raise SystemExit("queue is empty")
print(
f"received message_id={properties.message_id} "
f"redelivered={method.redelivered} body={body.decode()}"
)
channel.basic_ack(method.delivery_tag)
print("acked")Docker 路径使用 Compose 中的本地管理员和 vhost;URL 中的特殊字符必须进行百分号编码,不应把生产密码写入脚本:
export AMQP_URL='amqp://local_admin:replace-with-a-random-local-password@127.0.0.1:5672/local_lab'
python hello_rabbit.py publish
docker exec lab-rabbitmq rabbitmqctl list_queues -p local_lab \
name messages_ready messages_unacknowledged consumers
python hello_rabbit.py consume
docker exec lab-rabbitmq rabbitmqctl list_queues -p local_lab \
name messages_ready messages_unacknowledged consumers原生 Linux 路径使用刚创建的管理员和默认 vhost /,在 URL 中写成 %2F:
export AMQP_URL='amqp://local_admin:replace-with-the-entered-password@127.0.0.1:5672/%2F'
python hello_rabbit.py publish
sudo rabbitmqctl list_queues -p / \
name messages_ready messages_unacknowledged consumers
python hello_rabbit.py consume
sudo rabbitmqctl list_queues -p / \
name messages_ready messages_unacknowledged consumers发布后应看到 messages_ready=1,消费并 ack 后回到 0。若发布输出成功但 queue 不增加,检查 exchange、binding 和 routing key;若消费期间进程退出而没有 ack,消息会回到 ready 并在下次显示 redelivered=True。这条最小链只证明对象声明、发布、投递和 ack 能工作,不提供 publisher confirm、最小权限、业务幂等或节点容错。
3. 从第一条消息进入可维护拓扑
3.1 分开拓扑身份和运行身份
创建业务 vhost、拓扑账号和运行账号。add_user 不传密码时会交互输入,避免秘密进入 shell history:
docker exec lab-rabbitmq rabbitmqctl add_vhost orders
docker exec -it lab-rabbitmq rabbitmqctl add_user orders_topology
docker exec -it lab-rabbitmq rabbitmqctl add_user orders_app
docker exec lab-rabbitmq rabbitmqctl set_permissions -p orders orders_topology \
'^orders\..*' '^orders\..*' '^orders\..*'
docker exec lab-rabbitmq rabbitmqctl set_permissions -p orders orders_app \
'^$' '^orders\.events$' '^orders\.(created|dead)$'拓扑账号声明 exchange、queue 和 binding;运行账号只发布到 orders.events,并消费 orders.created 或 orders.dead。权限变更可能被现有 connection/channel 缓存,修改后让客户端重连再判断。
3.2 声明业务拓扑
创建 declare_topology.py:
import os
import pika
credentials = pika.PlainCredentials(
os.environ["RABBITMQ_USERNAME"],
os.environ["RABBITMQ_PASSWORD"],
)
params = pika.ConnectionParameters(
host=os.getenv("RABBITMQ_HOST", "127.0.0.1"),
port=int(os.getenv("RABBITMQ_PORT", "5672")),
virtual_host="orders",
credentials=credentials,
heartbeat=30,
)
with pika.BlockingConnection(params) as connection:
channel = connection.channel()
channel.exchange_declare(
exchange="orders.events", exchange_type="topic", durable=True
)
channel.exchange_declare(
exchange="orders.dlx", exchange_type="direct", durable=True
)
channel.queue_declare(
queue="orders.created",
durable=True,
arguments={"x-queue-type": "quorum"},
)
channel.queue_declare(
queue="orders.dead",
durable=True,
arguments={"x-queue-type": "quorum"},
)
channel.queue_bind(
queue="orders.created",
exchange="orders.events",
routing_key="commerce.order.created.v1",
)
channel.queue_bind(
queue="orders.dead",
exchange="orders.dlx",
routing_key="order.failed",
)使用拓扑身份运行:
export RABBITMQ_USERNAME='orders_topology'
export RABBITMQ_PASSWORD='replace-topology-password'
python declare_topology.py
docker exec lab-rabbitmq rabbitmqctl list_exchanges -p orders name type durable
docker exec lab-rabbitmq rabbitmqctl list_bindings -p orders \
source_name destination_name destination_kind routing_key
docker exec lab-rabbitmq rabbitmqctl list_queues -p orders name type durable3.3 用 Policy 管理可变规则
为主队列设置 dead letter、delivery limit 和 RabbitMQ 4.3 delayed retry。这里使用一条 queue 内部重试路线,不再额外声明一个从未使用的 TTL retry queue:
docker exec lab-rabbitmq rabbitmqctl set_policy -p orders orders-safety \
'^orders\.created$' \
'{"dead-letter-exchange":"orders.dlx","dead-letter-routing-key":"order.failed","dead-letter-strategy":"at-least-once","overflow":"reject-publish","max-length-bytes":1073741824,"delivery-limit":5,"delayed-retry-type":"failed","delayed-retry-min":30000,"delayed-retry-max":120000}' \
--apply-to quorum_queues --priority 10
docker exec lab-rabbitmq rabbitmqctl list_policies -p orders
docker exec lab-rabbitmq rabbitmqctl list_feature_flags name state \
| grep stream_queuefailed 只延迟会增加 delivery count 的失败,适合把连接中断或 basic.reject(requeue=true) 与普通业务返回分开;如果客户端使用 basic.nack(requeue=true),它既不增加 delivery count,也不会被这条 failed policy 延迟,必须改用相符的策略或由应用维护独立次数上限。示例中的 1 GiB 只是让本地实验拥有明确容量边界,生产值要由消息大小、积压窗口和恢复速率推导,并优先由 operator policy 施加。关键环境启用 at-least-once dead lettering 前,还要确认集群节点版本和 stream_queue feature flag;目标 DLX 或 queue 不可用时,死信继续作为源 queue 的 live message 占用该上限,达到上限后普通 publish 会被拒绝。
客户端声明参数写入对象自身,普通 policy 可以动态匹配资源,operator policy 则由运维施加上限。三者同时设置同一数值参数时,通常取更保守的有效值;但 queue type 这类不可变属性不能靠 policy 把已有对象原地转换。发布前应同时列出 queue arguments、普通 policy 和 operator policy,避免只看应用代码误判有效配置。
3.4 用 CLI 和 UI 识别日常状态
开发时先从对象关系和运行状态入手,不要一上来删除 queue:
docker exec lab-rabbitmq rabbitmqctl list_vhosts name
docker exec lab-rabbitmq rabbitmqctl list_users user tags
docker exec lab-rabbitmq rabbitmqctl list_permissions -p orders
docker exec lab-rabbitmq rabbitmqctl list_exchanges -p orders name type durable
docker exec lab-rabbitmq rabbitmqctl list_bindings -p orders \
source_name destination_name destination_kind routing_key
docker exec lab-rabbitmq rabbitmqctl list_queues -p orders \
name type durable arguments policy messages_ready \
messages_unacknowledged consumers
docker exec lab-rabbitmq rabbitmqctl list_connections \
user vhost peer_host peer_port state channels
docker exec lab-rabbitmq rabbitmqctl list_channels \
connection number consumer_count messages_unacknowledgedManagement UI 的 Connections、Channels、Exchanges、Queues and Streams、Admin 分别对应这些对象。页面图表便于学习最近状态,CLI 适合精确查询和自动化,HTTP API 适合集成;生产长期指标应交给 Prometheus,而不是依赖浏览器中有限的历史窗口。
4. 跑通 Confirm、Return 与订阅式消费
4.1 用异常而不是布尔值判断 Pika Confirm
创建 publish_once.py:
import json
import os
import pika
event_id = os.getenv("EVENT_ID", "evt-order-demo-1001-v1")
payload = {
"eventId": event_id,
"eventType": "OrderCreated",
"orderId": "demo-1001",
"schemaVersion": 1,
}
credentials = pika.PlainCredentials(
os.environ["RABBITMQ_USERNAME"],
os.environ["RABBITMQ_PASSWORD"],
)
params = pika.ConnectionParameters(
host=os.getenv("RABBITMQ_HOST", "127.0.0.1"),
port=int(os.getenv("RABBITMQ_PORT", "5672")),
virtual_host="orders",
credentials=credentials,
heartbeat=30,
blocked_connection_timeout=10,
)
try:
with pika.BlockingConnection(params) as connection:
channel = connection.channel()
channel.confirm_delivery()
channel.basic_publish(
exchange="orders.events",
routing_key=os.getenv(
"ROUTING_KEY", "commerce.order.created.v1"
),
body=json.dumps(payload).encode(),
mandatory=True,
properties=pika.BasicProperties(
content_type="application/json",
delivery_mode=pika.DeliveryMode.Persistent,
message_id=event_id,
type="commerce.order.created.v1",
),
)
except pika.exceptions.UnroutableError as exc:
raise SystemExit(f"publish returned as unroutable: {exc}") from exc
except pika.exceptions.NackError as exc:
raise SystemExit(f"broker nacked publish: {exc}") from exc
else:
print({"eventId": event_id, "brokerConfirmed": True})Pika 的 BlockingChannel.basic_publish() 在 confirm 模式成功时不返回 True;它在不可路由或 broker nack 时分别抛出 UnroutableError、NackError。没有异常返回,只能说明这次调用取得了相应 broker 确认,不能说明库存业务已经处理。
切换到运行身份并发布:
export RABBITMQ_USERNAME='orders_app'
export RABBITMQ_PASSWORD='replace-app-password'
python publish_once.py
docker exec lab-rabbitmq rabbitmqctl list_queues -p orders \
name messages_ready messages_unacknowledged consumers4.2 用 Mandatory Return 暴露错误路由
使用一个不存在 binding 的 routing key:
ROUTING_KEY='commerce.order.unknown.v1' \
EVENT_ID='evt-order-bad-route-v1' \
python publish_once.py脚本应以 publish returned as unroutable 退出,orders.created 不增加。若关闭 mandatory,该 publish 仍可能获得 confirm,却没有任何 queue 保存它。排查时依次看 exchange、binding 和 queue,而不是临时创建一个拼错名称的新队列:
docker exec lab-rabbitmq rabbitmqctl list_exchanges -p orders name type
docker exec lab-rabbitmq rabbitmqctl list_bindings -p orders \
source_name destination_name destination_kind routing_key
docker exec lab-rabbitmq rabbitmqctl list_queues -p orders \
name messages_ready messages_unacknowledged4.3 用 basic_consume 证明 Prefetch 与手动 Ack
创建 consume_once.py:
import json
import os
import pika
credentials = pika.PlainCredentials(
os.environ["RABBITMQ_USERNAME"],
os.environ["RABBITMQ_PASSWORD"],
)
params = pika.ConnectionParameters(
host=os.getenv("RABBITMQ_HOST", "127.0.0.1"),
port=int(os.getenv("RABBITMQ_PORT", "5672")),
virtual_host="orders",
credentials=credentials,
heartbeat=30,
)
with pika.BlockingConnection(params) as connection:
channel = connection.channel()
channel.basic_qos(prefetch_count=1)
def handle(ch, method, properties, body):
event = json.loads(body)
print({
"eventId": event["eventId"],
"deliveryTag": method.delivery_tag,
"redelivered": method.redelivered,
})
input("delivery is unacked; press Enter to ack: ")
ch.basic_ack(delivery_tag=method.delivery_tag)
ch.stop_consuming()
channel.basic_consume(
queue="orders.created", on_message_callback=handle, auto_ack=False
)
channel.start_consuming()在输入等待期间,另一个终端应看到 messages_unacknowledged=1。直接终止消费者进程,再次运行会看到 redelivered=True。这里使用 basic_consume,因为 QoS prefetch 对 basic.get 不生效。
5. 用 Outbox 与 Inbox 收敛两个不确定窗口
最小协议实验只证明 RabbitMQ 的状态变化。下面用 SQLite 建立一个可以在开发机运行的订单闭环;生产项目可把相同约束迁移到 PostgreSQL、MySQL 或其他事实库。
5.1 订单与 Outbox 同事务提交
创建 create_order.py:
import json
import sqlite3
event_id = "evt-order-demo-1001-v1"
order_id = "demo-1001"
payload = json.dumps({
"eventId": event_id,
"eventType": "OrderCreated",
"orderId": order_id,
"schemaVersion": 1,
})
with sqlite3.connect("orders.db") as db:
db.executescript("""
PRAGMA journal_mode=WAL;
CREATE TABLE IF NOT EXISTS orders (
order_id TEXT PRIMARY KEY,
status TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS outbox_event (
event_id TEXT PRIMARY KEY,
routing_key TEXT NOT NULL,
payload TEXT NOT NULL,
published_at TEXT
);
""")
db.execute("BEGIN IMMEDIATE")
db.execute(
"INSERT OR IGNORE INTO orders(order_id, status) VALUES (?, ?)",
(order_id, "CREATED"),
)
db.execute(
"""INSERT OR IGNORE INTO outbox_event
(event_id, routing_key, payload)
VALUES (?, ?, ?)""",
(event_id, "commerce.order.created.v1", payload),
)
db.commit()订单和待发布事件要么一起存在,要么一起回滚。此时尚未把 published_at 写成已完成。
5.2 只有 Confirm 和 Return 都通过才完成 Outbox
创建 relay_once.py:
import os
import sqlite3
import pika
db = sqlite3.connect("orders.db")
row = db.execute(
"""SELECT event_id, routing_key, payload
FROM outbox_event
WHERE published_at IS NULL
ORDER BY rowid LIMIT 1"""
).fetchone()
if row is None:
raise SystemExit("no pending outbox event")
event_id, routing_key, payload = row
credentials = pika.PlainCredentials(
os.environ["RABBITMQ_USERNAME"],
os.environ["RABBITMQ_PASSWORD"],
)
params = pika.ConnectionParameters(
host=os.getenv("RABBITMQ_HOST", "127.0.0.1"),
port=int(os.getenv("RABBITMQ_PORT", "5672")),
virtual_host="orders",
credentials=credentials,
heartbeat=30,
blocked_connection_timeout=10,
)
try:
with pika.BlockingConnection(params) as connection:
channel = connection.channel()
channel.confirm_delivery()
channel.basic_publish(
exchange="orders.events",
routing_key=routing_key,
body=payload.encode(),
mandatory=True,
properties=pika.BasicProperties(
content_type="application/json",
delivery_mode=pika.DeliveryMode.Persistent,
message_id=event_id,
type=routing_key,
),
)
except (pika.exceptions.UnroutableError, pika.exceptions.NackError):
raise
else:
db.execute(
"UPDATE outbox_event SET published_at=CURRENT_TIMESTAMP "
"WHERE event_id=? AND published_at IS NULL",
(event_id,),
)
db.commit()
finally:
db.close()python create_order.py
python relay_once.py
sqlite3 orders.db \
'select event_id, published_at from outbox_event;'连接在 confirm 返回前中断时,published_at 保持空值。Relay 稍后复用同一 event_id 发布,broker 内可能已有第一份消息,因此消费端必须幂等。生产系统还要用行锁或租约避免多个 relay 同时抢到同一 outbox,并限制未确认 publish 窗口。
5.3 数据库提交后再 Ack,重投时只产生一次业务结果
创建 inventory_consumer.py:
import json
import os
import sqlite3
import pika
inventory = sqlite3.connect("inventory.db")
inventory.executescript("""
PRAGMA journal_mode=WAL;
CREATE TABLE IF NOT EXISTS processed_event (
event_id TEXT PRIMARY KEY,
processed_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP
);
CREATE TABLE IF NOT EXISTS reservation (
order_id TEXT PRIMARY KEY,
status TEXT NOT NULL
);
""")
credentials = pika.PlainCredentials(
os.environ["RABBITMQ_USERNAME"],
os.environ["RABBITMQ_PASSWORD"],
)
params = pika.ConnectionParameters(
host=os.getenv("RABBITMQ_HOST", "127.0.0.1"),
port=int(os.getenv("RABBITMQ_PORT", "5672")),
virtual_host="orders",
credentials=credentials,
heartbeat=30,
)
with pika.BlockingConnection(params) as connection:
channel = connection.channel()
channel.basic_qos(prefetch_count=4)
def handle(ch, method, properties, body):
try:
event = json.loads(body)
event_id = event["eventId"]
order_id = event["orderId"]
except (json.JSONDecodeError, KeyError):
ch.basic_reject(method.delivery_tag, requeue=False)
return
try:
inventory.execute("BEGIN IMMEDIATE")
inventory.execute(
"INSERT INTO processed_event(event_id) VALUES (?)",
(event_id,),
)
inventory.execute(
"INSERT INTO reservation(order_id, status) VALUES (?, ?)",
(order_id, "RESERVED"),
)
inventory.commit()
print({"eventId": event_id, "businessApplied": True})
except sqlite3.IntegrityError:
inventory.rollback()
duplicate = inventory.execute(
"SELECT 1 FROM processed_event WHERE event_id=?",
(event_id,),
).fetchone()
if duplicate is None:
ch.basic_reject(method.delivery_tag, requeue=False)
return
print({"eventId": event_id, "duplicateIgnored": True})
if os.getenv("CRASH_AFTER_COMMIT") == "1":
os._exit(17)
ch.basic_ack(method.delivery_tag)
channel.basic_consume(
queue="orders.created", on_message_callback=handle, auto_ack=False
)
channel.start_consuming()先制造“数据库已提交、ack 尚未发送”的崩溃:
CRASH_AFTER_COMMIT=1 python inventory_consumer.py
python inventory_consumer.py
sqlite3 inventory.db 'select * from processed_event;'
sqlite3 inventory.db 'select * from reservation;'第二次消费应打印 duplicateIgnored,并对重投执行 ack。processed_event 和 reservation 各只有一行。这个实验解释了为什么重复投递可以接受,而重复业务副作用不能接受。
永久格式错误通过 basic_reject(requeue=False) 进入 orders.dead。局部瞬时错误若要使用当前 failed delayed retry,应执行 basic_reject(requeue=True);它会增加 delivery count,按线性退避再次投递,并在超过 delivery limit 后转入 DLX。basic_nack(requeue=True) 不增加 delivery count,不能依赖当前 policy 和 delivery limit 把它收敛。数据库整体不可用时应停止 consumer,不能让所有 delivery 一直重排。
6. 把同一语义接入 Spring AMQP
Python 实验把协议状态展示得最直接,真实 Java 项目还需要依赖、拓扑、序列化、发布确认、数据库事务、手动 ack 和停机行为连成完整链。下面以 Spring Boot 4.1.1、Java 21、Spring AMQP 4.1 和 PostgreSQL 为基线;已有项目应由自己的 Boot BOM 管理 Spring AMQP 与 RabbitMQ Java Client 版本,不要单独覆盖传递依赖制造不兼容组合。
6.1 建立项目依赖和真实配置
pom.xml 至少包含 AMQP、JDBC、PostgreSQL 和 Actuator:
<project xmlns="http://maven.apache.org/POM/4.0.0">
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-parent</artifactId>
<version>4.1.1</version>
</parent>
<groupId>demo</groupId>
<artifactId>orders-rabbitmq</artifactId>
<version>1.0.0</version>
<properties><java.version>21</java.version></properties>
<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-jdbc</artifactId>
</dependency>
<dependency>
<groupId>org.postgresql</groupId>
<artifactId>postgresql</artifactId>
<scope>runtime</scope>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-actuator</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
</dependencies>
<build>
<plugins>
<plugin>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-maven-plugin</artifactId>
</plugin>
</plugins>
</build>
</project>连接、confirm、return、manual ack 和 prefetch 使用框架真实属性,密码只从环境变量注入:
spring:
datasource:
url: ${JDBC_URL:jdbc:postgresql://127.0.0.1:5432/orders}
username: ${JDBC_USERNAME:orders_app}
password: ${JDBC_PASSWORD}
sql:
init:
mode: always
rabbitmq:
addresses: ${RABBITMQ_ADDRESSES:127.0.0.1:5672}
virtual-host: ${RABBITMQ_VHOST:orders}
username: ${RABBITMQ_USERNAME}
password: ${RABBITMQ_PASSWORD}
connection-timeout: 5s
requested-heartbeat: 30s
publisher-confirm-type: correlated
publisher-returns: true
dynamic: ${RABBITMQ_DYNAMIC:false}
template:
mandatory: true
listener:
simple:
acknowledge-mode: manual
prefetch: 16
concurrency: 2
max-concurrency: 8
default-requeue-rejected: false
lifecycle:
timeout-per-shutdown-phase: 30s
management:
endpoint:
health:
probes:
enabled: true
group:
readiness:
include: readinessState,db,rabbitdynamic=false 是运行身份的默认安全边界:应用会使用这些 bean 的名称消费和发布,但 RabbitAdmin 不会在每次启动时重新声明拓扑。部署阶段由 orders_topology 身份临时设置 RABBITMQ_DYNAMIC=true 创建对象,验证成功后再用只有 write/read 权限的运行身份启动;否则前文“拓扑身份与运行身份分离”会被启动期自动声明悄悄破坏。
创建 src/main/resources/schema.sql,让数据库唯一键承担重复边界:
CREATE TABLE IF NOT EXISTS processed_event (
event_id text PRIMARY KEY,
processed_at timestamptz NOT NULL DEFAULT now()
);
CREATE TABLE IF NOT EXISTS inventory_reservation (
order_id text PRIMARY KEY,
event_id text NOT NULL UNIQUE,
status text NOT NULL,
updated_at timestamptz NOT NULL DEFAULT now()
);
CREATE TABLE IF NOT EXISTS orders (
order_id text PRIMARY KEY,
status text NOT NULL,
created_at timestamptz NOT NULL DEFAULT now()
);
CREATE TABLE IF NOT EXISTS outbox_event (
event_id text PRIMARY KEY,
routing_key text NOT NULL,
payload jsonb NOT NULL,
created_at timestamptz NOT NULL DEFAULT now(),
published_at timestamptz
);6.2 声明拓扑并等待 Confirm 和 Return
拓扑对象由单独配置类声明。正式团队可以把声明权限交给部署任务;无论由谁创建,对象名称、类型和 arguments 都应来自同一份契约。
package demo.orders;
import org.springframework.amqp.core.*;
import org.springframework.amqp.support.converter.JacksonJsonMessageConverter;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import tools.jackson.databind.json.JsonMapper;
@Configuration
public class RabbitTopology {
public static final String EXCHANGE = "orders.events";
public static final String QUEUE = "orders.created";
public static final String ROUTING_KEY = "commerce.order.created.v1";
public static final String DLX = "orders.dlx";
public static final String DEAD_QUEUE = "orders.dead";
public static final String DEAD_ROUTING_KEY = "order.failed";
@Bean TopicExchange ordersExchange() {
return new TopicExchange(EXCHANGE, true, false);
}
@Bean Queue ordersCreated() {
return QueueBuilder.durable(QUEUE)
.quorum()
.build();
}
@Bean Binding ordersCreatedBinding(Queue ordersCreated,
TopicExchange ordersExchange) {
return BindingBuilder.bind(ordersCreated)
.to(ordersExchange).with(ROUTING_KEY);
}
@Bean DirectExchange ordersDlx() {
return new DirectExchange(DLX, true, false);
}
@Bean Queue ordersDead() {
return QueueBuilder.durable(DEAD_QUEUE).quorum().build();
}
@Bean Binding ordersDeadBinding(Queue ordersDead,
DirectExchange ordersDlx) {
return BindingBuilder.bind(ordersDead)
.to(ordersDlx).with(DEAD_ROUTING_KEY);
}
@Bean JacksonJsonMessageConverter jsonMessageConverter(JsonMapper json) {
return new JacksonJsonMessageConverter(json);
}
}Converter 复用 Spring Boot 自动配置的同一个 JsonMapper,因此 Outbox JSON、HTTP JSON 与 AMQP JSON 会共享自定义模块、命名策略和日期格式;另起一个 mapper 会在契约稍复杂时产生“数据库能读、消息却反序列化失败”的隐蔽分叉。
应用入口保持最小即可:
package demo.orders;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
@SpringBootApplication
public class OrdersApplication {
public static void main(String[] args) {
SpringApplication.run(OrdersApplication.class, args);
}
}事件对象保留稳定 ID。Outbox relay 调用发布器,而不是让订单事务直接跨数据库和 RabbitMQ 双写:
package demo.orders;
public record OrderCreated(
String eventId,
String orderId,
int schemaVersion
) {}订单接口不能直接“提交数据库后顺手发消息”。它只在一个数据库事务中写订单和 outbox;event ID 由业务键稳定推导,因此同一订单请求重试不会制造新事件:
package demo.orders;
import tools.jackson.core.JacksonException;
import tools.jackson.databind.json.JsonMapper;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
@Service
public class OrderService {
private final JdbcTemplate jdbc;
private final JsonMapper json;
public OrderService(JdbcTemplate jdbc, JsonMapper json) {
this.jdbc = jdbc;
this.json = json;
}
@Transactional
public String create(String orderId) throws JacksonException {
String eventId = "evt-order-" + orderId + "-v1";
OrderCreated event = new OrderCreated(eventId, orderId, 1);
jdbc.update(
"INSERT INTO orders(order_id,status) VALUES (?,?) " +
"ON CONFLICT DO NOTHING",
orderId, "CREATED");
jdbc.update(
"INSERT INTO outbox_event(event_id,routing_key,payload) " +
"VALUES (?,?,?::jsonb) ON CONFLICT DO NOTHING",
eventId, RabbitTopology.ROUTING_KEY,
json.writeValueAsString(event));
return eventId;
}
}下面的 relay 一次取一行,等待 confirm/return 后才标记完成,足以跑通单实例教学闭环。生产多实例应增加 FOR UPDATE SKIP LOCKED、租约或状态抢占,限制 in-flight 数,并定期扫描超龄 pending;即使两个 relay 偶尔重复发布,也必须沿用原 event ID:
package demo.orders;
import tools.jackson.databind.json.JsonMapper;
import java.util.List;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.stereotype.Component;
@Component
public class OutboxRelay {
private record Pending(
String eventId, String routingKey, String payload) {}
private final JdbcTemplate jdbc;
private final JsonMapper json;
private final OrderPublisher publisher;
public OutboxRelay(JdbcTemplate jdbc, JsonMapper json,
OrderPublisher publisher) {
this.jdbc = jdbc;
this.json = json;
this.publisher = publisher;
}
public boolean relayOne() throws Exception {
List<Pending> rows = jdbc.query(
"SELECT event_id,routing_key,payload::text FROM outbox_event " +
"WHERE published_at IS NULL ORDER BY created_at LIMIT 1",
(rs, rowNum) -> new Pending(
rs.getString("event_id"), rs.getString("routing_key"),
rs.getString("payload")));
if (rows.isEmpty()) return false;
Pending row = rows.get(0);
publisher.publish(
row.routingKey(),
json.readValue(row.payload(), OrderCreated.class));
return jdbc.update(
"UPDATE outbox_event SET published_at=now() " +
"WHERE event_id=? AND published_at IS NULL",
row.eventId()) == 1;
}
}为实验提供两个明确触发点。/internal/outbox/relay-once 只是教学入口,生产必须放在受控内部网络并由调度器或后台 worker 驱动:
package demo.orders;
import java.util.Map;
import org.springframework.web.bind.annotation.*;
@RestController
public class OrderController {
private final OrderService orders;
private final OutboxRelay relay;
public OrderController(OrderService orders, OutboxRelay relay) {
this.orders = orders;
this.relay = relay;
}
@PostMapping("/orders/{orderId}")
Map<String, Object> create(@PathVariable String orderId) throws Exception {
return Map.of("eventId", orders.create(orderId), "outbox", "pending");
}
@PostMapping("/internal/outbox/relay-once")
Map<String, Object> relayOne() throws Exception {
return Map.of("published", relay.relayOne());
}
}package demo.orders;
import java.util.concurrent.TimeUnit;
import org.springframework.amqp.AmqpException;
import org.springframework.amqp.rabbit.connection.CorrelationData;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.stereotype.Component;
@Component
public class OrderPublisher {
private final RabbitTemplate rabbit;
public OrderPublisher(RabbitTemplate rabbit) {
this.rabbit = rabbit;
}
public void publish(String routingKey, OrderCreated event) throws Exception {
CorrelationData correlation = new CorrelationData(event.eventId());
rabbit.convertAndSend(
RabbitTopology.EXCHANGE,
routingKey,
event,
message -> {
message.getMessageProperties().setMessageId(event.eventId());
message.getMessageProperties().setType(routingKey);
return message;
},
correlation
);
CorrelationData.Confirm confirm = correlation.getFuture()
.get(10, TimeUnit.SECONDS);
if (!confirm.ack()) {
throw new AmqpException("broker nack: " + confirm.reason());
}
if (correlation.getReturned() != null) {
throw new AmqpException(
"unroutable: " + correlation.getReturned().getReplyText());
}
}
}CorrelationData 会把 return 写入同一发布上下文,并在 confirm future 完成前提供它。只有 confirm 为 ack 且 return 为空,relay 才把 outbox 标记为 published。Future 超时、连接中断或进程退出仍是未知结果,outbox 保持 pending,并用原 event ID 再发;不能创建新 ID 绕过消费者幂等。
6.3 数据库提交后再 Ack
把事务方法放在独立 Bean,使 @Transactional 代理真正生效。首次 event ID 插入成功才更新业务表;重复事件直接返回 false:
package demo.orders;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
@Service
public class InventoryService {
private final JdbcTemplate jdbc;
public InventoryService(JdbcTemplate jdbc) {
this.jdbc = jdbc;
}
@Transactional
public boolean reserve(OrderCreated event) {
int first = jdbc.update(
"INSERT INTO processed_event(event_id) VALUES (?) " +
"ON CONFLICT DO NOTHING",
event.eventId());
if (first == 0) return false;
jdbc.update(
"INSERT INTO inventory_reservation(order_id,event_id,status) " +
"VALUES (?,?,?) ON CONFLICT(order_id) DO UPDATE SET " +
"event_id=excluded.event_id,status=excluded.status,updated_at=now()",
event.orderId(), event.eventId(), "RESERVED");
return true;
}
}Listener 等 reserve() 返回时,独立数据库事务已经提交,随后才 ack。局部瞬时数据库异常用 basicReject(tag, true) 接入前文 failed delayed retry;永久业务错误用 false 进入 DLX。数据库整体离线时应停止或暂停容器,不应让所有 delivery 独立重试。
package demo.orders;
import com.rabbitmq.client.Channel;
import java.io.IOException;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.dao.TransientDataAccessException;
import org.springframework.stereotype.Component;
@Component
public class InventoryListener {
private final InventoryService inventory;
public InventoryListener(InventoryService inventory) {
this.inventory = inventory;
}
@RabbitListener(queues = RabbitTopology.QUEUE)
public void handle(OrderCreated event, Message message, Channel channel)
throws IOException {
long tag = message.getMessageProperties().getDeliveryTag();
try {
boolean changed = inventory.reserve(event);
if (Boolean.getBoolean("demo.crash-after-commit")) {
System.err.println("test hook: halt after database commit");
Runtime.getRuntime().halt(17);
}
channel.basicAck(tag, false);
System.out.printf("eventId=%s businessChanged=%s%n",
event.eventId(), changed);
} catch (TransientDataAccessException retryable) {
channel.basicReject(tag, true);
} catch (IllegalArgumentException permanent) {
channel.basicReject(tag, false);
} catch (RuntimeException unknown) {
System.err.printf("eventId=%s unclassified=%s%n",
event.eventId(), unknown.getClass().getName());
channel.basicReject(tag, true);
}
}
}Ack 本身失败会带来重投,因此 inbox 唯一键仍不可省略。示例把方法体内无法分类的运行时异常显式 basicReject(tag, true),使其进入有界 delayed retry;JSON 转换等发生在 listener 方法之前的致命异常由容器处理,default-requeue-rejected=false 会拒绝且不重新入主队列,并按 DLX 设置转移。两类路径都要监控并最终收敛,不能让默认行为靠猜。
6.4 验证启动、重投和优雅停机
先让 PostgreSQL、RabbitMQ 拓扑和凭据就绪,再启动应用:
export JDBC_PASSWORD='replace-database-password'
export JDBC_USERNAME='orders_app'
export JDBC_URL='jdbc:postgresql://127.0.0.1:5432/orders'
export PSQL_URL='postgresql://orders_app@127.0.0.1:5432/orders'
export RABBITMQ_USERNAME='orders_topology'
export RABBITMQ_PASSWORD='replace-topology-password'
RABBITMQ_DYNAMIC=true \
SPRING_RABBITMQ_LISTENER_SIMPLE_AUTO_STARTUP=false \
mvn spring-boot:start
curl --fail http://127.0.0.1:8080/actuator/health/readiness
mvn spring-boot:stop拓扑创建完成后停止部署进程。本地完整实验改用同时拥有最小 write/read 权限的 orders_app 启动:
export RABBITMQ_USERNAME='orders_app'
export RABBITMQ_PASSWORD='replace-app-password'
mvn spring-boot:start
curl --fail http://127.0.0.1:8080/actuator/health/readiness生产部署再把 relay 与 consumer 拆成不同进程或 profile:前者只启用 OrderPublisher 并使用 orders_publisher,后者只启用 InventoryListener 并使用 orders_consumer。一个 JVM 只有一组 RabbitMQ 连接凭据时,不能声称发布与消费已经最小权限分离;若暂时同进程,就应明确记录这是过渡期的组合身份。
Readiness 必须显示 readinessState、db 和 rabbit 都为 UP。接着创建订单并明确触发一次 relay:
curl --fail -X POST http://127.0.0.1:8080/orders/demo-1001
curl --fail -X POST \
http://127.0.0.1:8080/internal/outbox/relay-once
psql "$PSQL_URL" -c \
"select event_id,published_at from outbox_event order by created_at desc limit 5;"
mvn spring-boot:stopspring-boot:start 在后台启动应用,因此同一终端可以继续执行 curl;PSQL_URL 使用 libpq 能识别的 postgresql://,不是应用的 jdbc:postgresql://,命令会交互询问数据库密码。第一条响应应显示 outbox 为 pending,第二条显示 published=true,随后 outbox 获得 published_at,listener 写入一份 processed_event 和一份 reservation。
再复现“数据库已提交、ack 前进程退出”。终端 A 用只允许隔离实验的 JVM 属性启动前台进程:
JAVA_TOOL_OPTIONS='-Ddemo.crash-after-commit=true' \
mvn spring-boot:run终端 B 创建另一笔订单并触发 relay:
curl --fail -X POST http://127.0.0.1:8080/orders/crash-1002
curl -X POST http://127.0.0.1:8080/internal/outbox/relay-onceConsumer 在事务提交后会让整个 JVM 以 17 退出,第二个 curl 可能因进程退出而断开,这正是未知结果。回到终端 B,以正常模式重启并再次扫描 pending outbox:
mvn spring-boot:start
curl --fail -X POST \
http://127.0.0.1:8080/internal/outbox/relay-once
psql "$PSQL_URL" -c \
"select event_id,count(*) from processed_event where event_id='evt-order-crash-1002-v1' group by event_id;"
psql "$PSQL_URL" -c \
"select order_id,count(*) from inventory_reservation where order_id='crash-1002' group by order_id;"
mvn spring-boot:stop两个 count 都应为 1;日志至少一次显示重复被忽略。测试钩子必须只保留在教学或测试 profile,生产构建应删除。若直接调用 OrderPublisher 而没有先写 outbox,这个演练就失去了生产端未知窗口,不能算完整链路。
停机时先从入口摘除实例,等待 listener container 停止接收新 delivery,并给已在处理的数据库事务留下完成时间。超过 shutdown phase 仍未完成的 delivery 会随 channel 关闭而重投;因此优雅停机减少重复,但不能代替幂等。自动恢复后还要重新验证 channel、confirm、consumer 和拓扑属性,connection 恢复本身不代表应用已经可服务。
7. 把开发节点迁移到生产安全基线
7.1 关闭明文入口并启用 mTLS
证书由团队 PKI 签发,服务端证书 SAN 必须包含客户端实际连接的 FQDN。CA、服务端证书和私钥安装到只允许 rabbitmq 读取的目录:
sudo install -d -o rabbitmq -g rabbitmq -m 0750 /etc/rabbitmq/tls
sudo install -o rabbitmq -g rabbitmq -m 0644 ca_certificate.pem \
/etc/rabbitmq/tls/ca_certificate.pem
sudo install -o rabbitmq -g rabbitmq -m 0644 server_certificate.pem \
/etc/rabbitmq/tls/server_certificate.pem
sudo install -o rabbitmq -g rabbitmq -m 0640 server_key.pem \
/etc/rabbitmq/tls/server_key.pem写入 /etc/rabbitmq/rabbitmq.conf。示例要求客户端证书,同时仍使用 RabbitMQ 用户/vhost 权限完成 AMQP 认证授权;Management 和 Prometheus 先绑定管理网回环地址:
listeners.tcp = none
listeners.ssl.default = 5671
ssl_options.cacertfile = /etc/rabbitmq/tls/ca_certificate.pem
ssl_options.certfile = /etc/rabbitmq/tls/server_certificate.pem
ssl_options.keyfile = /etc/rabbitmq/tls/server_key.pem
ssl_options.verify = verify_peer
ssl_options.fail_if_no_peer_cert = true
ssl_options.versions.1 = tlsv1.3
ssl_options.versions.2 = tlsv1.2
management.tcp.ip = 127.0.0.1
management.tcp.port = 15672
prometheus.tcp.ip = 127.0.0.1
prometheus.tcp.port = 15692重启前先保留现有管理会话和配置副本,确认应用已经具备 AMQPS 参数;否则关闭 5672 会造成全部客户端同时离线。重启后检查实际环境、TLS 版本、listener 和日志:
sudo systemctl restart rabbitmq-server
sudo rabbitmq-diagnostics --silent tls_versions
sudo rabbitmq-diagnostics listeners
sudo rabbitmq-diagnostics environment
sudo journalctl -u rabbitmq-server -n 100 --no-pagerlistener 应出现 5671 的 amqp/ssl,不再出现普通 5672。先验证证书链、SNI/主机名和客户端证书,再用真实 AMQP 客户端验证用户、vhost、权限和消息收发:
openssl s_client \
-connect rmq01.example.com:5671 \
-servername rmq01.example.com \
-CAfile ca_certificate.pem \
-cert client_certificate.pem \
-key client_key.pem \
-verify_return_error \
-verify_hostname rmq01.example.com去掉 -cert/-key 的反向测试应被服务端拒绝。openssl 成功仍没有执行 AMQP 登录;Pika 客户端还要加载证书、保持主机名校验并提供最小权限账号:
import ssl
import pika
context = ssl.create_default_context(cafile="ca_certificate.pem")
context.load_cert_chain("client_certificate.pem", "client_key.pem")
context.check_hostname = True
context.verify_mode = ssl.CERT_REQUIRED
params = pika.ConnectionParameters(
host="rmq01.example.com",
port=5671,
virtual_host="orders",
credentials=pika.PlainCredentials("orders_app", "secret-from-vault"),
ssl_options=pika.SSLOptions(context, "rmq01.example.com"),
heartbeat=30,
)
with pika.BlockingConnection(params) as connection:
print("AMQPS and vhost authentication succeeded")7.2 分离管理员、拓扑、发布、消费和监控身份
身份按动作拆分,不按团队人数拆分:
| 身份 | configure | write | read | 说明 |
|---|---|---|---|---|
| 集群管理员 | 受控全局 | 受控全局 | 受控全局 | 只用于运维变更,不进入应用 |
orders_topology | ^orders\..* | ^orders\..* | ^orders\..* | 部署阶段声明和迁移拓扑 |
orders_publisher | ^$ | ^orders\.events$ | ^$ | 只能向订单 exchange 发布 |
orders_consumer | ^$ | ^$ | ^orders\.created$ | 只能消费主队列 |
| 监控用户 | ^$ | ^$ | ^$ | 用 monitoring tag 或受控 Prometheus 网络读取指标 |
创建用户时交互输入密码,随后逐个设置权限;不要让一个 orders_app 长期同时发布、消费和读取 dead queue:
sudo rabbitmqctl add_user orders_publisher
sudo rabbitmqctl add_user orders_consumer
sudo rabbitmqctl set_permissions -p orders orders_publisher \
'^$' '^orders\.events$' '^$'
sudo rabbitmqctl set_permissions -p orders orders_consumer \
'^$' '^$' '^orders\.created$'
sudo rabbitmqctl list_permissions -p orders密码、客户端私钥和 cookie 是三种不同 secret。应用密码和私钥来自密钥系统,证书轮换先让客户端信任新旧 CA/证书,再换服务端,最后移除旧信任;Erlang cookie 只在节点与受控 CLI 之间一致,权限必须为仅节点账号可读。Management、EPMD、节点分布式端口和 Prometheus 不应与业务 AMQP 端口共享公网暴露面。
8. 从单节点进入三节点 Quorum Queue
8.1 建立同版本集群
本地多数派实验使用独立 Compose 文件,三个节点共享同一 Erlang cookie,但各自拥有数据卷。生产节点还应分布到不同故障域,并使用 TLS、DNS、资源限制和受控 cookie 文件;本地三容器不代表跨可用区灾难恢复。
services:
rmq1:
image: rabbitmq:4.3.5-management
hostname: rmq1
environment:
RABBITMQ_ERLANG_COOKIE: replace-with-one-random-cluster-cookie
RABBITMQ_NODENAME: rabbit@rmq1
ports:
- "127.0.0.1:5673:5672"
- "127.0.0.1:15673:15672"
volumes:
- rmq1-data:/var/lib/rabbitmq
rmq2:
image: rabbitmq:4.3.5-management
hostname: rmq2
environment:
RABBITMQ_ERLANG_COOKIE: replace-with-one-random-cluster-cookie
RABBITMQ_NODENAME: rabbit@rmq2
volumes:
- rmq2-data:/var/lib/rabbitmq
rmq3:
image: rabbitmq:4.3.5-management
hostname: rmq3
environment:
RABBITMQ_ERLANG_COOKIE: replace-with-one-random-cluster-cookie
RABBITMQ_NODENAME: rabbit@rmq3
volumes:
- rmq3-data:/var/lib/rabbitmq
volumes:
rmq1-data:
rmq2-data:
rmq3-data:启动节点,等 rmq1 就绪后再加入其余成员:
docker compose -f compose.cluster.yaml up -d
docker compose -f compose.cluster.yaml exec rmq1 rabbitmqctl await_startup
docker compose -f compose.cluster.yaml exec rmq2 rabbitmqctl await_startup
docker compose -f compose.cluster.yaml exec rmq3 rabbitmqctl await_startup
docker compose -f compose.cluster.yaml exec rmq2 rabbitmqctl stop_app
docker compose -f compose.cluster.yaml exec rmq2 rabbitmqctl reset
docker compose -f compose.cluster.yaml exec rmq2 rabbitmqctl join_cluster rabbit@rmq1
docker compose -f compose.cluster.yaml exec rmq2 rabbitmqctl start_app
docker compose -f compose.cluster.yaml exec rmq3 rabbitmqctl stop_app
docker compose -f compose.cluster.yaml exec rmq3 rabbitmqctl reset
docker compose -f compose.cluster.yaml exec rmq3 rabbitmqctl join_cluster rabbit@rmq1
docker compose -f compose.cluster.yaml exec rmq3 rabbitmqctl start_app
docker compose -f compose.cluster.yaml exec rmq1 \
rabbitmq-diagnostics cluster_status这个集群拥有全新的数据目录,单节点实验中的 vhost、用户和 queue 不会自动出现。先在 rmq1 创建身份与权限,再通过宿主机的 5673 复用前文拓扑脚本:
docker compose -f compose.cluster.yaml exec rmq1 \
rabbitmqctl add_vhost orders
docker compose -f compose.cluster.yaml exec -it rmq1 \
rabbitmqctl add_user orders_topology
docker compose -f compose.cluster.yaml exec -it rmq1 \
rabbitmqctl add_user orders_app
docker compose -f compose.cluster.yaml exec rmq1 \
rabbitmqctl set_permissions -p orders orders_topology \
'^orders\..*' '^orders\..*' '^orders\..*'
docker compose -f compose.cluster.yaml exec rmq1 \
rabbitmqctl set_permissions -p orders orders_app \
'^$' '^orders\.events$' '^orders\.(created|dead)$'
RABBITMQ_HOST=127.0.0.1 \
RABBITMQ_PORT=5673 \
RABBITMQ_USERNAME=orders_topology \
RABBITMQ_PASSWORD='replace-topology-password' \
python declare_topology.py新声明的 quorum queue 默认会在可用集群成员中形成复制组;实际成员数仍以 quorum_status 为准。若只有一个成员,先停止故障实验并检查集群成员、默认 quorum 初始组大小和声明日志,不能把单副本 queue 当成三副本:
docker compose -f compose.cluster.yaml exec rmq1 \
rabbitmq-queues quorum_status --vhost orders orders.created8.2 分别验证失去一个成员与失去多数派
保持 publisher 使用稳定 event ID。先停止一个 follower:
docker compose -f compose.cluster.yaml stop rmq3三成员队列仍有两票,多数派存在,publish 可以在重新选主或副本状态稳定后继续获得 confirm。再停止第二个成员:
docker compose -f compose.cluster.yaml stop rmq2此时只剩一个成员,quorum queue 不能安全推进。Publisher 应得到 nack、超时或连接层错误,outbox 保持 pending;不能为了恢复写入而强行重建同名 queue。恢复成员后重新查询 quorum 状态,再用原 event ID 重发:
docker compose -f compose.cluster.yaml start rmq2 rmq3
docker compose -f compose.cluster.yaml exec rmq1 \
rabbitmq-queues quorum_status --vhost orders orders.created生产网络分区还要检查 RabbitMQ 4.3 分区行为。RabbitMQ 4.3 的元数据存储和 quorum queue 都基于 Raft;“节点进程仍在”不能代替 leader 和多数派状态。
8.3 让节点加入、退出和替换都保持多数派
生产节点先固定可解析的长节点名。在每台主机的 /etc/rabbitmq/rabbitmq-env.conf 写入自己的 FQDN,三台机器只共用 cookie,不共用数据目录:
NODENAME=rabbit@rmq01.example.com
USE_LONGNAME=truermq02、rmq03 分别改成自己的名字。三台机器必须能双向解析 FQDN,节点间只在集群网络开放 EPMD 4369 和分布式通信 25672。把同一份随机 cookie 安装为 /var/lib/rabbitmq/.erlang.cookie,属主设为 rabbitmq:rabbitmq、权限设为 0400;不要把 cookie 当作应用密码,也不要提交到配置仓库。
只有全新、无业务数据的节点才执行 reset 后加入。以下命令在新节点 rmq02 上运行,目标是已存在的 rmq01:
sudo rabbitmqctl stop_app
sudo rabbitmqctl reset
sudo rabbitmqctl join_cluster rabbit@rmq01.example.com --longnames
sudo rabbitmqctl start_app
sudo rabbitmq-diagnostics cluster_status --longnames现有成员绝不能为了“重新加入”随手 reset;reset 会清除本节点的集群元数据和本地数据。节点进入 cluster 也只加入共享元数据,已有 quorum queue 不会自动把新节点纳入复制组。先针对指定队列增加成员并检查结果:
sudo rabbitmq-queues add_member \
--vhost orders orders.created rabbit@rmq04.example.com
sudo rabbitmq-queues quorum_status \
--vhost orders orders.created节点较多时可以按受控模式把现有 quorum queue 和 stream 扩到新节点,再重平衡 leader。先用测试 vhost 验证模式,避免一次性搬迁全部队列:
sudo rabbitmq-queues grow rabbit@rmq04.example.com all \
--vhost-pattern '^orders$' --queue-pattern '^orders\.'
sudo rabbitmq-queues rebalance all下线节点不能只停止 systemd。先确认移除后每个 quorum queue 仍保有多数派,再把该节点上的复制成员迁出:
sudo rabbitmq-diagnostics check_if_node_is_quorum_critical
sudo rabbitmq-queues shrink rabbit@rmq03.example.com
sudo rabbitmq-queues quorum_status --vhost orders orders.created如果只迁移单个队列,可用 delete_member;成员变化本身也需要现有多数派同意。永久失联且无法恢复的节点,必须先逐队列证明剩余成员仍有多数派,完成副本补齐后才从 cluster metadata 忘记该节点。替换节点的安全顺序是:加入新节点 → grow/add_member → 等待同步 → 检查 quorum critical → shrink/delete_member 旧节点 → 再移除旧成员。每一步都以 queue 的真实成员与 leader 为准,不以 cluster_status 里的节点数量代替。
9. 观察容量、流控与积压恢复
队列、消费者和节点要一起观察:
docker exec lab-rabbitmq rabbitmqctl list_queues -p orders \
name type messages_ready messages_unacknowledged message_bytes \
consumers consumer_capacity
docker exec lab-rabbitmq rabbitmq-diagnostics check_local_alarms
docker exec lab-rabbitmq rabbitmq-diagnostics memory_breakdown
docker exec lab-rabbitmq rabbitmqctl statusrabbitmqctl status 包含磁盘余量和水位信息;rabbitmq-diagnostics disk_space 不是可用的 RabbitMQ 4.3 命令。节点触发 memory 或 disk alarm 时会阻塞发布连接,应用表现为 confirm 延迟、blocked callback 或超时。此时先降低到达率、恢复消费和释放受控磁盘空间,不要先提高水位掩盖耗尽。
Prefetch 可从“单消费者目标吞吐 × 单条处理时间”估算起点。例如每秒处理 40 条、P95 为 0.25s,可从 10 开始压测。最终值要同时看 consumer capacity、unacked、进程内存、业务延迟和消费者退出后的重投峰值。
监控至少包含:publish/confirm/return rate、ready/unacked、deliver/ack/redelivery rate、consumer capacity、连接和 channel 数、blocked connection、quorum leader/member、dead letter 增长、磁盘和内存 alarm。应用侧再补 outbox pending age、inbox duplicate、处理耗时和业务失败分类。
在每个节点启用 Prometheus 插件,而不是只抓一个随机节点:
sudo rabbitmq-plugins enable rabbitmq_prometheus
sudo rabbitmq-diagnostics listeners
curl --fail http://127.0.0.1:15692/metrics | head默认入口是 15692。生产环境应绑定受控监控网或回环地址,由采集器逐节点抓取;节点级、集群级和对象级指标的聚合方式不同,不能把同一个 cluster counter 从三台节点简单相加。Management API 适合人工诊断,Prometheus 指标适合持续采集,两者都不能暴露到公网。
由于前文把 15692 绑定在 loopback,默认使用每节点一个 Prometheus Agent 本地抓取,再通过 remote_write 送往中心平台。每台节点都使用 127.0.0.1:15692,并把自己的 RabbitMQ node name 写入 label;下面是 rmq01 的片段:
scrape_configs:
- job_name: rabbitmq
scrape_interval: 15s
static_configs:
- targets:
- 127.0.0.1:15692
labels:
rabbitmq_node: rabbit@rmq01.example.com
remote_write:
- url: https://metrics.example.com/api/v1/writermq02、rmq03 分别替换 label,remote write 凭据从 secret 文件注入。若不部署本机 Agent,才把 prometheus.tcp.ip 改为每台主机的受控监控网 IP,并用主机防火墙只允许中心 Prometheus 来源;不能一边绑定 127.0.0.1,一边让中心服务直接抓 FQDN,也不能为了省事监听公网地址。
第一轮告警不要只写“队列大于某个数”,而要把信号连到责任:ready 连续多个窗口净增长时通知消费负责人并计算排空时间;任一节点出现 memory/disk alarm 时立即限制发布并通知平台负责人;confirm P95/P99 与 outbox oldest age 同时上升时定位 broker 接管窗口;关键 quorum queue 在线成员少于预期或 leader 缺失时冻结节点变更。具体指标名以当前 /metrics 实际输出和 RabbitMQ Grafana dashboard 为准,规则上线前要用停消费者、填充受控磁盘和停止一个 follower 的演练证明它能触发并恢复。
容量先从业务窗口反推,而不是从磁盘剩余百分比猜:
净排空速率 = ack rate - publish rate
积压恢复时间 = ready backlog / 净排空速率
逻辑积压字节 ≈ publish rate × 平均编码后消息字节 × 最长积压秒数
物理磁盘预算 ≈ 逻辑积压字节 × quorum 副本数 × 存储开销与安全余量若每秒进入 2,000 条、平均 4 KiB、允许积压 30 分钟,逻辑 payload 已约 13.7 GiB;三副本 quorum queue 还要乘副本数,并留出日志段、索引、合并、重平衡和恢复余量。恢复后消费者只能比生产者多处理 500/s 时,360 万 条积压至少还需 2 小时排空。只看当前 ready 数、不看净排空速率,就无法回答何时恢复。
连接、channel、consumer 和文件描述符也要预算。连接用于网络生命周期,channel 用于协议并发,不能让每个请求都创建连接;连接池泄漏会先吃掉 socket 和文件描述符。生产节点通常至少把文件描述符上限提高到 65536,但最终值要由连接数、队列文件和主机约束验证。systemd 环境可用受控 override 设置:
[Service]
LimitNOFILE=65536执行 sudo systemctl daemon-reload && sudo systemctl restart rabbitmq-server 后,用 rabbitmqctl status 对照操作系统限制确认生效。阈值不应照抄本文数字:用代表性消息大小、confirm 模式、持久化类型、消费者耗时和节点故障场景压测,分别记录稳态、单副本恢复、积压排空和磁盘逼近水位时的吞吐与 P95/P99。
告警必须能落到动作:ready 持续增长且 ack 低于 publish,先找慢消费者或下游;unacked 与处理延迟一起升高,检查 prefetch、线程阻塞和数据库;confirm 延迟与 blocked connection 同时上升,检查 memory/disk alarm;leader 缺失或成员不足,停止变更并恢复多数派;outbox oldest age 增长但 broker 指标正常,检查 relay;dead queue 增长先按错误类型采样。恢复判定则看趋势反转、积压预计时间、confirm/ack 回稳和业务唯一状态,而不是只把告警手工关闭。
10. 备份、恢复、升级与清理
10.1 Definitions 不包含 Queue 中的消息
RabbitMQ definitions 可以导出 vhost、用户、权限、exchange、queue、binding 和 policy,但不会导出 queue 中的消息。导出文件包含安全与拓扑信息,应按敏感配置保存:
docker exec lab-rabbitmq rabbitmqctl export_definitions \
/tmp/rabbitmq-definitions.json
docker cp lab-rabbitmq:/tmp/rabbitmq-definitions.json \
./rabbitmq-definitions.json
docker exec lab-rabbitmq rm -f /tmp/rabbitmq-definitions.json
chmod 600 ./rabbitmq-definitions.json在不支持 POSIX 权限的工作站上,应使用操作系统原生 ACL 把文件限制给当前账号,不能把 chmod 的结果当作跨平台权限证明。Definitions 含用户密码哈希、内部名称和权限关系,不应提交到仓库或进入普通构建日志。
启动一个独立空白容器,导入后逐层检查对象:
docker run -d --name lab-rabbitmq-restore \
-p 127.0.0.1:5674:5672 \
-p 127.0.0.1:15674:15672 \
rabbitmq:4.3.5-management
docker exec lab-rabbitmq-restore rabbitmqctl await_startup
docker cp ./rabbitmq-definitions.json \
lab-rabbitmq-restore:/tmp/rabbitmq-definitions.json
docker exec lab-rabbitmq-restore rabbitmqctl import_definitions \
/tmp/rabbitmq-definitions.json
docker exec lab-rabbitmq-restore rm -f /tmp/rabbitmq-definitions.json
docker exec lab-rabbitmq-restore rabbitmqctl list_vhosts name
docker exec lab-rabbitmq-restore rabbitmqctl list_user_permissions orders_app
docker exec lab-rabbitmq-restore rabbitmqctl list_exchanges -p orders name type
docker exec lab-rabbitmq-restore rabbitmqctl list_bindings -p orders \
source_name destination_name destination_kind routing_key
docker exec lab-rabbitmq-restore rabbitmqctl list_policies -p orders
docker exec lab-rabbitmq-restore rabbitmqctl list_queues -p orders name type导入 definitions 后 queue 是空的,因为消息从未包含在文件中。使用恢复端口和原 event ID 重发,再运行幂等消费者:
RABBITMQ_PORT=5674 \
RABBITMQ_USERNAME=orders_app \
RABBITMQ_PASSWORD='replace-app-password' \
EVENT_ID='evt-order-demo-1001-v1' \
python publish_once.py
RABBITMQ_PORT=5674 \
RABBITMQ_USERNAME=orders_app \
RABBITMQ_PASSWORD='replace-app-password' \
python inventory_consumer.py现有 inventory.db 已处理过该 event ID 时,消费者应打印 duplicateIgnored 并 ack;使用独立空白业务库时应只建立一份库存预留。完成隔离测试后执行 docker rm -f -v lab-rabbitmq-restore,definitions 文件则按备份保留策略保存或安全删除。
直接复制正在运行节点的数据目录不是一致备份;quorum queue 也不能只拿一个成员的数据目录冒充完整副本。生产消息恢复通常来自业务事实库、outbox、上游日志或受控的跨集群复制。
恢复结束至少要确认:原 event ID 能再次发布,业务 inbox 只产生一份状态,ready/unacked 能回落,dead letter 目标存在,quorum queue 有 leader 和多数派。Definitions 成功导入但消息来源缺失,仍不算业务恢复。
10.2 升级按节点和 Feature Flag 推进
下面以 RabbitMQ 4.2.x → 4.3.5、Erlang 27.x 为例。滚动升级只支持相邻 minor,不允许从更老 minor 跨级跳到 4.3;插件也必须存在与目标版本兼容的构建。先在隔离环境恢复 definitions 和代表性消息,完成发布、消费、重投与恢复演练,再检查集群:
sudo rabbitmqctl -q --formatter pretty_table \
list_feature_flags name state stability provided_by
sudo rabbitmqctl enable_feature_flag all
sudo rabbitmq-diagnostics cluster_status
sudo rabbitmq-diagnostics check_local_alarms
sudo rabbitmqctl list_queues name type state messages_ready \
messages_unacknowledged这里的 enable_feature_flag all 在仍是完整 4.2 集群时执行,用来启用旧版本已经提供的 stable flags。存在 disabled stable flag、分区、资源 alarm、quorum-critical 节点、未同步成员、插件不兼容或目标包不可获得时停止升级;不要带病进入混合版本。
一次只处理一个非 quorum-critical 节点。升级该节点前,让工具确认除它以外仍有在线多数派,并把客户端连接与 leader 迁走:
sudo rabbitmq-diagnostics check_if_node_is_quorum_critical
sudo rabbitmq-upgrade await_online_quorum_plus_one
sudo rabbitmq-upgrade drain然后停止服务,用已锁定的软件源把 RabbitMQ 升到 4.3.5,Erlang 保持兼容的 27.x,再启动并逐层验证:
sudo systemctl stop rabbitmq-server
sudo apt-get install rabbitmq-server=4.3.5-1
sudo systemctl start rabbitmq-server
sudo rabbitmq-diagnostics ping
sudo rabbitmqctl status
sudo rabbitmq-diagnostics cluster_status
sudo rabbitmq-diagnostics listeners
sudo journalctl -u rabbitmq-server -n 100 --no-pager包版本后缀以仓库实际值为准,先用 apt-cache madison rabbitmq-server 确认,不能盲抄 -1。等待所有 quorum queue/stream 成员追平,真实应用重新获得 confirm、consume、ack,ready/unacked 和业务错误率稳定后,才能处理下一节点。若在真正停止服务前放弃该节点升级,可执行 sudo rabbitmq-upgrade revive 让被 drain 的节点恢复承担连接和 leader。
全部节点都升级并稳定后,再评估启用 4.3 新 feature flags:
sudo rabbitmqctl enable_feature_flag all
sudo rabbitmqctl -q --formatter pretty_table \
list_feature_flags name state stability provided_by
sudo rabbitmq-queues rebalance allFeature flag 启用后通常不能关闭,RabbitMQ 也不支持把数据节点原地降级回旧版本。因此“包降级”不是可靠 rollback。需要真正版本回退能力时,应提前建设蓝绿集群、兼容双写或受控 Federation/Shovel,验证新集群后切流;一旦在原集群启用新 flag,就把恢复策略视为向前修复或切回仍保留的旧集群,而不是覆盖数据目录。
10.3 清理只针对本次实验资源
停止消费者并确认没有 unacked 后,使用拓扑身份删除本次对象,再撤销用户和 vhost。共享实例禁止 purge 整个未知 queue,也禁止删除 volume。
个人隔离环境可以保留数据停止:
docker compose down只有确认 Compose 项目名、volume 名和其中没有其他项目或排障现场后,才在该隔离目录执行:
docker compose down -vPython 虚拟环境、SQLite 数据库和 definitions 文件可能仍含凭证、payload 或业务样例;按团队的数据处理要求删除或归档,不要默认提交到仓库。
四、问题处理
1. 先定位消息停在哪一层
RabbitMQ 故障可以沿同一条链收敛:
DNS/TCP/TLS
→ authentication/vhost/authorization
→ exchange 与 mandatory return
→ binding 与 queue
→ ready / unacked
→ ack / reject / delayed retry / DLX
→ quorum leader 与多数派
→ memory / disk alarm
→ outbox / inbox / 业务状态先保留 event ID、routing key、异常类型、queue 状态和业务幂等记录。不要用重启、purge 或批量补发代替分层判断。
先建立一份最小现场,后续每个分支都从它继续,而不是每查一步就改一项配置:
date -Is
sudo rabbitmqctl version
sudo rabbitmq-diagnostics ping
sudo rabbitmq-diagnostics listeners
sudo rabbitmq-diagnostics cluster_status
sudo rabbitmq-diagnostics check_local_alarms
sudo rabbitmqctl list_queues -p orders \
name type state messages_ready messages_unacknowledged consumers
sudo journalctl -u rabbitmq-server --since '-15 min' --no-pager容器环境把 sudo ... 换成 docker exec <container> ...,并同时保存 docker inspect、容器日志和 Compose 项目名。记录命令时间、节点名、版本和目标 vhost;否则几分钟后的自动恢复会让“当时为什么失败”无法复盘。变更前写清预期输出和回退动作,变更后用相同命令重跑,只有基础设施指标与业务结果都恢复才结束。
2. 安装、启动与节点身份故障
2.1 包已安装,但服务无法启动
现象与影响: systemctl status 显示 failed,rabbitmq-diagnostics ping 无法连接,本机没有 AMQP listener。此时应用重试只会放大日志,消息是否丢失取决于上游有没有 outbox 或重放源。
先看包组合、服务错误和有效配置:
apt-cache policy rabbitmq-server erlang-base erlang-ssl
sudo systemctl status rabbitmq-server --no-pager
sudo journalctl -u rabbitmq-server -b -n 200 --no-pager
sudo rabbitmq-diagnostics environment
sudo ss -lntp | grep -E ':(4369|5671|5672|15672|25672)\b'日志若明确指出 Erlang 版本不兼容,回到官方兼容矩阵,用同一受控仓库安装匹配组合;不能只替换 RabbitMQ 包后继续启动。若是 rabbitmq.conf 解析错误,恢复上一份已验证配置,逐项加入新配置并重启验证。若端口被其他进程占用,先识别进程归属,再修改冲突服务或规划 listener;不要直接结束未知进程。
修复后依次要求 systemctl is-active 返回 active、ping 成功、listeners 出现预期协议、日志没有新的 boot error,最后完成一次真实 publish/consume。预防措施是固定 RabbitMQ/Erlang 软件源与版本、在预发布节点执行配置检查和冷启动、保留上一版本包与配置,但不要把复制活动数据目录当成回滚。
2.2 CLI 说节点不存在或 cookie 不匹配
现象与影响: 服务进程可能存在,但 rabbitmqctl 报 nodedown、invalid challenge reply 或连接到错误节点;集群成员还可能因 FQDN、短名混用而互相不可见。先核对服务实际节点名、主机解析、CLI 身份和 cookie 来源:
sudo rabbitmq-diagnostics status
sudo rabbitmq-diagnostics erlang_cookie_sources
hostname -f
getent hosts rmq01.example.com rmq02.example.com rmq03.example.com
sudo rabbitmq-diagnostics cluster_status --longnames如果服务使用长名,CLI 也要带 --longnames;如果只是当前 shell 用户读取了另一份 cookie,优先用 sudo -u rabbitmq 执行 CLI,而不是覆盖运行节点 cookie。集群节点之间的 cookie 确实不一致时,先停止准备修复的节点,备份其配置并确认其不是 quorum-critical,再按受控流程安装正确 cookie。绝不能为了让 CLI 连上就把一个未知集群的 cookie 覆盖到运行中的数据节点。
修复后验证节点名只出现一次、三台机器双向解析一致、cluster_status 的 running nodes 与预期相同,再逐队列确认 leader 和成员。预防措施是把节点 FQDN、USE_LONGNAME、cookie secret 的来源和文件权限纳入主机配置管理,并在节点加入前运行 DNS/cookie 预检。
2.3 容器反复重启或“改了环境变量却没变化”
先查看退出原因、健康检查、挂载和真实 volume:
docker compose ps
docker compose logs --tail 200 rabbitmq
docker inspect lab-rabbitmq --format \
'{{.State.Status}} {{.State.ExitCode}} {{json .Mounts}}'
docker volume lsRabbitMQ 首次启动后,默认用户、definitions、cookie 和节点身份已经写入持久化状态;修改初始化环境变量不会重写既有数据。共享环境应通过 rabbitmqctl、definitions 迁移或配置管理显式变更;只有确认 Compose 项目、volume 绝对归属和无保留数据的个人实验,才可按清理章节重建。恢复后不仅看容器为 running,还要验证 listener、用户/vhost、queue 类型和一条真实消息。
3. 连接、认证或权限失败
3.1 Management UI 可用,AMQP 仍失败
先确认应用连接的是 5672 或 5671,不是 15672。再看 broker listener、vhost、用户权限和认证日志:
docker exec lab-rabbitmq rabbitmq-diagnostics listeners
docker exec lab-rabbitmq rabbitmqctl list_vhosts name
docker exec lab-rabbitmq rabbitmqctl list_users user tags
docker exec lab-rabbitmq rabbitmqctl list_user_permissions orders_app
docker logs --tail 200 lab-rabbitmq若日志完全没有应用连接记录,问题仍在地址、端口、DNS、防火墙或 TLS 之前;若出现 PLAIN login refused,检查用户名/密码;若出现 access to vhost refused,检查 vhost 是否存在及用户权限;若连接成功但声明/发布/消费失败,再进入资源权限。密码不应作为 authenticate_user 参数进入命令历史;认证结果用应用身份的最小 AMQP 连接来确认。
默认 guest 只允许 loopback。远程应用应创建独立账号,不应通过关闭 loopback 保护放开默认口令。
3.2 ACCESS_REFUSED 或 TLS 握手失败
ACCESS_REFUSED 要区分 authentication、vhost access 和资源正则:
sudo rabbitmqctl list_permissions -p orders
sudo rabbitmqctl list_user_permissions orders_publisher
sudo rabbitmqctl list_user_permissions orders_consumer
sudo journalctl -u rabbitmq-server --since '-10 min' --no-pager只有 configure 失败时,不要顺手授予 .* / .* / .*;确认运行进程是否错误启用了动态声明。write 失败就核对 exchange 正则,read 失败就核对 queue 正则。修复后以对应最小身份分别执行声明、发布和消费反向测试,确保 publisher 仍不能读 queue、consumer 仍不能写 exchange。
TLS 失败先用 openssl s_client 保留握手结果,再检查 CA、服务器 SAN、SNI/连接主机名、证书有效期、客户端证书链和协议版本:
openssl s_client -connect rmq01.example.com:5671 \
-servername rmq01.example.com -CAfile ca_certificate.pem \
-cert client_certificate.pem -key client_key.pem \
-verify_return_error -verify_hostname rmq01.example.comVerify return code: 0 只证明 TLS 链和主机名通过,仍需 AMQP 登录、vhost 权限和消息收发。临时关闭证书校验只能帮助定位,不能进入应用配置;最终还要做一次无客户端证书的反向测试,确认 mTLS 确实拒绝它。
4. 拓扑与发布端故障
4.1 Confirm 超时
Confirm 超时首先是未知结果,不是明确失败。停止扩大未确认窗口,保留 event ID 和 outbox pending;随后检查 connection blocked、节点 alarm、目标 queue 类型、quorum leader/member、磁盘延迟和网络:
docker exec lab-rabbitmq rabbitmq-diagnostics check_local_alarms
docker exec lab-rabbitmq rabbitmqctl list_connections \
name state channels send_pend
docker exec lab-rabbitmq rabbitmq-queues quorum_status \
--vhost orders orders.createdconnection blocked 或 alarm 存在时先降低发布流量、恢复消费和资源,不要继续提高 confirm 并发;queue 无 leader 时先恢复多数派;网络抖动但 broker 可能已经接收时,结果保持 unknown。恢复后用同一 event ID 重发,让 consumer inbox 收敛可能的重复。新建 event ID 会绕开幂等约束。验证标准是 confirm 延迟回到基线、outbox oldest age 开始下降、同一 event ID 在业务库仍只有一个结果。
4.2 Confirm 成功但 Queue 为零
先看 return handler 是否收到不可路由消息,再核对 exchange 类型、routing key 和 binding:
sudo rabbitmqctl list_exchanges -p orders name type durable
sudo rabbitmqctl list_bindings -p orders \
source_name destination_name destination_kind routing_key arguments
sudo rabbitmqctl list_queues -p orders name type messages_readyConfirm 只说明 broker 处理了 publish;没有 queue 匹配时也可能确认。mandatory return 指向 binding/routing key 错误,完全没有 return 则检查客户端是否真的启用 mandatory 和 return callback。修复 binding 后用一个新的测试 event ID 发布,要求 confirm=ack、return 为空、目标 queue 增加 1,再由消费者处理;不要为“试一下”删除共享 exchange。
4.3 PRECONDITION_FAILED
同名 queue 或 exchange 已经用不同 durable、exclusive、type 或 argument 声明。错误日志会指出不等价的字段;同时查询 definitions 或 Management API 中的当前对象。durable/exclusive/type 这类声明时属性不能用 policy 覆盖,创建版本化新对象、双写或迁移消费者后再退役旧对象;delivery limit、dead-letter 等适合 policy 的可变项则改 policy。不要让应用在启动时反复删除共享 queue。验证时重启两个不同版本的应用,二者都不应再争抢声明。
5. 消费端与业务一致性故障
5.1 Ready 增长但没有处理
查看 consumers、consumer capacity、deliver/ack rate 和应用连接:
sudo rabbitmqctl list_queues -p orders \
name messages_ready messages_unacknowledged consumers consumer_capacity
sudo rabbitmqctl list_consumers -p orders
sudo rabbitmqctl list_connections name user state channelsconsumers=0 时修复应用启动、queue 名和 read 权限;consumer 存在但 capacity 长期接近零时,抓线程栈并检查数据库锁、外部请求和线程池;capacity 高但 ready 仍涨,说明总处理能力低于到达率,需要先限流再扩消费或优化单条耗时。修复后要看到 ack rate 持续高于 publish rate,并用公式估算 ready 清零时间。
5.2 Unacked 长期增长
Unacked 高说明消息已经离开 ready。先找出它集中在哪个 connection/channel,而不是立即重启全部消费者:
sudo rabbitmqctl list_queues -p orders \
name messages_ready messages_unacknowledged consumers
sudo rabbitmqctl list_channels connection number consumer_count \
prefetch_count messages_unacknowledged
sudo rabbitmqctl list_consumers -p orders若几乎全部 unacked 集中在一个 channel,抓对应实例线程栈、数据库锁等待和外部请求耗时,检查异常路径是否漏掉 ack/reject;若均匀分布且单条处理都慢,问题是下游或总容量;若 unacked 约等于 consumer × prefetch 且 ready 也增长,prefetch 只是把积压搬进应用。先摘除坏实例或暂停入口,再缩小 prefetch/并发窗口并修复阻塞点。强制断开连接会让该连接全部 unacked 重投,执行前必须证明 event ID 和 inbox 唯一键能承受峰值。恢复后要求 unacked 回到稳定区间、ack rate 高于 publish rate、业务延迟回落;预防措施是按实例记录处理耗时、in-flight 和停机排空时间。
5.3 unknown delivery tag
Broker 通常会以 PRECONDITION_FAILED - unknown delivery tag 关闭当前 channel,随后该 channel 的未确认消息重投。先把 broker 日志时间与应用 stack trace 对齐,并检查 channel 数是否反复创建:
sudo journalctl -u rabbitmq-server --since '-10 min' --no-pager
sudo rabbitmqctl list_connections name user state channels
sudo rabbitmqctl list_channels connection number messages_unacknowledged同一 tag 被 ack 两次,说明成功路径与 finally/回调重复确认;tag 来自另一个 channel,说明线程池丢失了 channel 归属;重连后旧 tag 继续出现,说明缓存了协议状态。修复为“delivery、处理回调、ack 都绑定收到它的 channel”,不在线程间共享非线程安全 channel,连接恢复时丢弃所有旧 tag。用一次正常 ack、一次 ack 前进程退出的实验复测:前者不再关闭 channel,后者只重投且业务结果仍唯一;持续告警 channel exception rate 防止回归。
5.4 数据库已提交但消息重复
先用应用日志关联 redelivered、message ID、event ID 和订单号,再查 inbox 与业务表:
SELECT event_id, count(*) FROM processed_event
GROUP BY event_id HAVING count(*) > 1;
SELECT order_id, event_id, status FROM inventory_reservation
WHERE order_id = 'demo-1001';inbox 已有 event ID 且业务状态正确,说明是 ack 前崩溃或 confirm 未知带来的正常至少一次重投,第二次只 ack;inbox 没有记录但业务已变化,说明副作用和去重不在同一事务;同一业务键出现两个不同 event ID,则生产者重试错误地生成了新 ID。先暂停有副作用的 consumer,按审计记录补偿或建立唯一约束,再恢复受控重放,不能用 purge 掩盖重复。复测必须主动在提交后、ack 前退出进程,并确认数据库仍只有一份状态;预防措施是稳定业务幂等键、同事务 inbox 和重复计数指标。
6. Retry 与 Dead Letter 故障
6.1 消息形成热循环
持续 basic.nack(requeue=true) 不增加 delivery count,也不受当前 failed policy 和 delivery limit 约束,会让消息立即或按其他 returned 策略再次可用,形成无上限循环。先看 redelivery/ack rate、CPU、同一 event ID 日志,以及消息的 x-acquired-count 与 x-delivery-count;前者表示被分配次数,后者才表示失败投递次数。x-acquired-count 快速增加而 x-delivery-count 不变,正是普通 nack return 在循环。
立即暂停问题 consumer 或把该消息隔离,避免它挤占整个 prefetch。局部瞬时错误改用 basic.reject(requeue=true) 配合 failed delayed retry;系统级依赖失败则暂停 consumer;永久错误 reject false 进入 DLX。修复后用一个故意失败的 event ID 观察 30/60/90/120 秒线性间隔、超过 delivery limit 后进入 dead queue,并确认健康消息仍能前进。预防措施是代码审查 ack/reject 分类、按 event ID 限制日志和对 redelivery/ack 比率告警。
6.2 消息没有进入 Dead Queue
按“源 policy → 生效 policy → DLX → binding → 目标 queue”查询:
sudo rabbitmqctl list_policies -p orders
sudo rabbitmqctl list_queues -p orders name type policy arguments state
sudo rabbitmqctl list_exchanges -p orders name type durable
sudo rabbitmqctl list_bindings -p orders \
source_name destination_name destination_kind routing_key
sudo rabbitmqctl list_queues -p orders \
name messages_ready messages_unacknowledged源 queue 没有匹配 policy 时修正 pattern/apply-to/priority;policy 已生效但 DLX 不存在时先恢复 exchange;DLX 存在但目标不增时修正 routing key/binding;目标 queue 无 leader时先恢复多数派。At-least-once dead lettering 下,目标不可用会让源 queue 保留消息并重试;默认策略则可能丢失。修复后发布一条明确永久失败的测试消息,要求源 queue 减少、dead queue 增加、原 event ID 和失败头仍在;不要手工复制 payload 后删除原消息。预防措施是把 DLX 目标纳入拓扑部署和恢复演练,而不是等出错才创建。
6.3 Dead Queue 快速增长
先记录 dead queue 增长率、最早消息年龄和主 queue 状态,再用只读检查消费者抽样少量消息,保持 requeue,不在 UI 中批量 ack。按 event type、schema version、异常类和生产者版本聚类:单一 schema 暴增通常是契约回归,多个类型同时暴增更像数据库/TLS等共享依赖,少量固定 event ID 反复出现则是毒消息或幂等缺口。
永久错误先修数据或部署兼容消费者;瞬时故障先恢复依赖;未经分类不得直接批量重放。建立隔离重放 worker,保持原 event ID、限制每秒速率,并同时监控主 queue、dead queue、confirm、数据库锁和业务错误。先重放 1 条,再 10 条,再逐步扩量;每一批都要求主流程成功、inbox 不重复、dead queue 净下降。预防措施是保存失败原因/原 routing key、按 schema 建兼容窗口,并把 dead queue oldest age 与增长率设为产品告警。
7. 资源、流控与容量故障
7.1 Memory、Disk Alarm 或文件描述符耗尽
现象与影响: publisher connection 被 blocked、confirm 延迟陡增、节点拒绝新连接,日志出现 memory/disk alarm 或 too many open files。先保存增速最大的 queue、连接/channel 数、磁盘和内存分解:
sudo rabbitmq-diagnostics check_local_alarms
sudo rabbitmq-diagnostics memory_breakdown
sudo rabbitmqctl status
sudo rabbitmqctl list_queues name message_bytes messages_ready \
messages_unacknowledged --sort message_bytes
sudo rabbitmqctl list_connections name user state channels send_penddisk alarm 时先限制 producer、恢复或扩容受控存储、确认日志与未知文件归属;memory alarm 时检查 ready/unacked、连接/channel 泄漏和消息体大小;FD 耗尽时定位连接风暴或文件预算,并核对 systemd 与进程实际 limit。直接调高 watermark、删除未知 queue、purge 或清空数据目录会把容量事故变成数据事故。
解除 alarm 后仍要确认 blocked connection 自动恢复、confirm 延迟回落、quorum 副本追平、ready 以可预测速率下降。预防措施是同时告警绝对余量与消耗速率,对连接/channel 设预算,用故障压测校准磁盘和内存,而不是把水位当日常容量目标。
7.2 连接或 Channel 数持续增长
按 user、peer host 和 state 查连接,再从应用连接池指标定位实例:
sudo rabbitmqctl list_connections name user peer_host peer_port state channels
sudo rabbitmqctl list_channels connection name number consumer_count \
messages_unacknowledged连接随请求数单调增长通常是客户端未复用/未关闭;单连接 channel 无界增长则是 channel 生命周期泄漏。先摘除有问题的应用实例或限流,不能粗暴关闭全部生产连接。修复客户端后观察连接/channel 回到稳定平台,并在发布超时、listener 重启和自动恢复场景重复验证。
8. Quorum Queue 与集群成员故障
8.1 Quorum Queue 没有 Leader 或多数派
rabbitmq-diagnostics cluster_status
rabbitmq-queues quorum_status --vhost orders orders.created命令应在安装了 RabbitMQ CLI 的节点或对应容器内执行。若成员只是暂时离线,优先恢复原主机、DNS、cookie 和网络;若永久损坏,先确认剩余成员是否仍有多数派,再按“新增成员—同步—移除旧成员”的流程替换。没有多数派时不要 force delete 同名 queue,也不要重建空 queue 接管流量;publisher outbox 保持 pending、consumer 停止推进是预期的一致性结果。恢复后必须看到 leader、预期成员数和在线多数派,并用原 event ID 收敛未知发布。
8.2 节点需要下线,却是 quorum-critical
rabbitmq-diagnostics check_if_node_is_quorum_critical 非零退出意味着现在停止该节点会让至少一个 quorum queue 或 stream 失去可用多数派。用 quorum_status 找出受影响队列,恢复其他成员或先把新成员加入复制组;等待同步完成后再检查。不能因为 cluster 仍有三台机器就假设所有 queue 都是三副本,也不能把“维护窗口已到”当成绕过多数派的理由。
8.3 Leader 或副本分布严重倾斜
节点加入或故障恢复后,leader 和成员可能集中在少数节点。先看 Management 的 queue leader/member 分布、各节点 CPU/磁盘/网络,再限定 vhost 和 queue pattern 预演范围:
sudo rabbitmq-diagnostics cluster_status
sudo rabbitmq-queues rebalance quorum \
--vhost-pattern '^orders$' --queue-pattern '^orders\.'只在无 alarm、网络稳定、quorum 成员在线且积压没有快速增长的低峰执行;否则选主和数据移动会放大故障。若只是 leader 倾斜,rebalance 即可;若新节点根本不是 queue member,先 grow/add_member 并等待同步,rebalance 不能凭空增加副本。完成后逐个 quorum_status,要求关键 queue 成员数正确、leader 分布改善、confirm/ack 没有异常峰值;把加入节点后的 grow 与 rebalance 固化进扩容流程。
9. Definitions、恢复与升级故障
9.1 Definitions 导入成功,消息却不见了
这是预期边界,不是导入工具漏数据。Definitions 只恢复 vhost、用户、权限、exchange、queue、binding 和 policy,不包含 queue payload。先核对恢复目标是否为空白隔离环境,再检查拓扑和消息重建来源:
sudo rabbitmqctl list_vhosts name
sudo rabbitmqctl list_bindings -p orders \
source_name destination_name destination_kind routing_key
sudo rabbitmqctl list_queues -p orders name type messages_ready若业务以数据库/outbox 为事实源,用原 event ID 受控重放;若没有任何重建来源,应把它记录为备份设计缺口,而不是伪造“恢复完成”。恢复完成至少应确认拓扑一致、代表性消息可发布消费、inbox 不重复、quorum 成员正常和积压可恢复。
9.2 Definitions 导入出现对象冲突
目标环境已有同名但不同 type、durable 或 arguments 的对象时,导入可能失败或留下部分对象。先保存 import 返回和 broker 日志,再在目标列出实际对象与 policy:
sudo rabbitmqctl list_exchanges -p orders name type durable
sudo rabbitmqctl list_queues -p orders name type durable arguments policy
sudo rabbitmqctl list_bindings -p orders \
source_name destination_name destination_kind routing_key
sudo rabbitmqctl list_policies -p orders若 type/durable 等不可变属性不同,不能靠提高 policy 优先级修复;为新对象使用版本化名称,建立双路由、排空旧 queue 再退役。若只是可变 policy 不同,修正 definitions 中 policy 来源后在隔离节点重导。不要直接覆盖生产,也不要在部分导入后手工删除未知对象。修复后的演练必须从空白节点重复一次,并完成代表性 publish/consume、权限反向测试和 queue type 检查,证明过程可重复而不是靠手工点选补齐。
9.3 滚动升级中节点无法重新加入
立即暂停后续节点,保住仍在线的多数派。检查目标节点的 RabbitMQ/Erlang 实际版本、插件、feature flags、节点名/cookie和启动日志:
sudo rabbitmqctl version
erl -noshell -eval 'io:format("~s~n", [erlang:system_info(otp_release)]), halt().'
sudo rabbitmq-plugins list -e
sudo rabbitmqctl -q --formatter pretty_table \
list_feature_flags name state stability provided_by
sudo journalctl -u rabbitmq-server -b -n 200 --no-pager若只是升级前 drain 后决定中止、服务尚未替换,可 rabbitmq-upgrade revive;若新包已经修改数据或启用了新 feature flag,不要尝试原地降级。修复兼容包/插件并让该节点重新在线、queue 成员追平、业务 confirm/ack 稳定后才能继续。需要旧版本业务回退时切到预先保留的蓝绿集群;没有蓝绿方案就按向前修复处理并明确风险。
10. 恢复后回到业务结果
服务进程变绿只完成基础设施恢复。回到最初的 evt-order-demo-1001-v1,publish 最终应得到 confirm 且没有 mandatory return;orders.created 应重新拥有 leader 和多数派,ready/unacked 恢复稳定;inbox 中只能有一个 event ID,库存预留也只能有一份。Delayed retry 和 dead queue 不再持续增长、outbox pending age 回落,并且未知 publish 已用原 event ID 收敛之后,再撤销临时诊断账号、放宽权限和测试资源。这条业务结果完整恢复,才说明 RabbitMQ、客户端和数据库之间的责任重新接上。
