RabbitMQ 部署、可靠投递与队列治理工具手册
从“发布成功”却没有订单说起
订单服务打印了“消息发送成功”,库存服务却没有收到消息。RabbitMQ 管理台能打开,连接数也不为零,队列深度甚至一直是 0。如果只检查端口和进程,这个现场看起来完全正常;真正的问题可能是消息发到了不存在的 exchange、routing key 没匹配任何 binding、生产者没有处理 mandatory return,或者消费者在业务事务提交前就自动 ack 了消息。
RabbitMQ 可靠性不是一个开关,而是三次责任转移:
生产者把消息交给 exchange,并等待 broker 的 publisher confirm。exchange 根据类型、routing key 和 binding 把消息路由到一个或多个 queue。消费者完成业务副作用后发送 ack,broker 才能删除该 delivery。
Confirm 与 consumer ack 彼此独立。Confirm 只说明 broker 已按当前队列类型接管消息,不说明消费者处理成功;ack 只说明消费者放弃再次投递的权利,不证明数据库事务、外部 API 或后续事件一定成功。RabbitMQ 的可靠性指南也把数据安全定义为 broker、publisher 和 consumer 的共同责任。
这条链路决定了排障顺序:先看发布有没有被确认,再看消息有没有路由进队列,然后区分 ready、unacked 和死信,最后才检查业务处理。如果不先分层,重启 broker、purge 队列或者反复补发只会扩大重复与丢失。
先建立 exchange、queue 与 binding 的心智模型
AMQP 0-9-1 客户端不是把消息直接发送给普通队列,而是发布到 exchange。Exchange 不存消息,它根据 exchange type 和 routing key 计算目标;queue 才保存等待投递的消息;binding 是 exchange 到 queue 的路由规则。
四种常用 exchange type 解决的是不同路由问题:
direct 要求 routing key 精确匹配,适合命令和确定的任务类型。topic 用 * 匹配一个词、# 匹配零个或多个词,适合有稳定分类维度的领域事件。fanout 忽略 routing key,把消息复制到全部绑定队列,适合广播刷新和多订阅者通知。
headers 根据消息头匹配,表达力强但治理和排障成本更高,通常不作为团队默认方案。
同一条消息可以进入多个队列,因此“发布一条”不等于“只处理一次”。同一个队列上的多个消费者是竞争消费,一条 delivery 同一时刻只交给其中一个消费者;多个独立队列则各自保存副本。需要每个下游都处理时,应为每个订阅者建立独立队列,而不是让多个业务共用一个队列名。
RabbitMQ 还提供空名称的 default exchange。每个队列会以自己的队列名自动绑定到它,便于快速验证,却会把路由设计隐藏在队列名里。正式项目宜显式声明 exchange、queue 和 binding,让拓扑可以审查、迁移和授权。
选择部署形态时先判断故障责任
单机二进制安装最适合学习 Erlang/RabbitMQ 服务目录、日志和系统服务,但容易污染开发机,并把 Erlang 与 RabbitMQ 的版本兼容交给每位开发者。单容器启动最短,适合一次性协议实验;容器删除后状态是否保留完全取决于 volume。Compose 可以把端口、持久化、健康检查和资源命名纳入仓库,是项目开发依赖的常用基线。
共享开发实例降低了每个人的资源消耗,却把 vhost、权限、资源前缀、TTL 和清理责任变成必选项。托管服务把节点、磁盘和升级交给供应商,但客户端仍然要正确处理 confirm、ack、重连、幂等和凭证轮换。Kubernetes Operator 适合平台团队声明式交付集群,不能消除多数派、磁盘水位和队列热点问题。
生产环境常见三种拓扑:
| 拓扑 | 组成与工作方式 | 适用信号 | 主要故障模式 |
|---|---|---|---|
| 单节点 | 一个 broker,classic queue 本地保存 | 可重建任务、低成本内部系统 | 节点或磁盘故障即中断,未复制消息可能不可用 |
| 三节点集群 + quorum queue | 元数据集群化,关键队列使用 Raft 多数派复制 | 订单、支付指令等不能接受单节点丢失 | 失去多数派后队列拒绝继续推进;磁盘和网络延迟放大尾延迟 |
| 托管多可用区 | 云服务管理节点和故障替换,客户端经服务端点接入 | 团队不希望自建值守体系 | 配额、跨区流量、能力差异和供应商故障仍需业务降级 |
RabbitMQ 集群共享用户、vhost、exchange、binding 等元数据,但 queue 内容是否复制取决于队列类型。官方集群指南明确建议集群运行在可靠 LAN 内;跨 WAN 连接应考虑 Federation 或 Shovel,而不是把一个集群横跨机房。
在开发机启动可清理的 RabbitMQ
RabbitMQ 4.3 是当前受支持发布线,补丁版本应由团队在升级评审后固定;不要使用 latest。官方发布信息用于确认支持状态,Docker Official Image提供 rabbitmq:4.3.2-management 标签,其中 5672 是 AMQP 端口,15672 是 management HTTP 端口。
先创建 .env,真实密码只保存在本机或密钥系统:
RABBITMQ_IMAGE=rabbitmq:4.3.2-management
RABBITMQ_AMQP_PORT=5672
RABBITMQ_UI_PORT=15672
RABBITMQ_DEFAULT_USER=te_admin
RABBITMQ_DEFAULT_PASS=replace-with-a-local-random-password
RABBITMQ_DEFAULT_VHOST=te_dev然后创建 compose.yaml:
services:
rabbitmq:
image: ${RABBITMQ_IMAGE:?set RABBITMQ_IMAGE}
container_name: te-rabbitmq
hostname: te-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:-te_dev}
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,避免开发机把管理台和 AMQP 裸露到局域网。启动后不要只看容器状态:
docker compose up -d rabbitmq
docker compose ps rabbitmq
docker exec te-rabbitmq rabbitmq-diagnostics -q ping
docker exec te-rabbitmq rabbitmq-diagnostics status
docker exec te-rabbitmq rabbitmqctl list_vhosts name tracing
docker exec te-rabbitmq rabbitmqctl list_users user tagsping 成功只证明 Erlang 节点可响应 CLI;应用还需要通过 AMQP 认证、进入目标 vhost 并获得资源权限。管理台可通过 http://127.0.0.1:15672 观察开发实例,但 HTTP API 不应替代应用的 AMQP 发布消费链。
环境变量创建默认用户只发生在空节点首次启动。若 volume 已经保存 RabbitMQ 数据,修改 .env 不会重置旧用户和密码。这不是 Compose 缓存,而是防止重启时覆盖现存安全状态的设计。
把资源声明和应用身份拆开
开发管理员负责创建 vhost、用户和拓扑;运行时应用只获得自己需要的资源正则权限。先创建两个身份:
docker exec te-rabbitmq rabbitmqctl add_vhost te_orders
docker exec -it te-rabbitmq rabbitmqctl add_user orders_topology
docker exec -it te-rabbitmq rabbitmqctl add_user orders_app
docker exec te-rabbitmq rabbitmqctl set_user_tags orders_topology management
docker exec te-rabbitmq rabbitmqctl set_permissions -p te_orders orders_topology \
"^orders\\..*" "^orders\\..*" "^orders\\..*"
docker exec te-rabbitmq rabbitmqctl set_permissions -p te_orders orders_app \
"^$" "^orders\\.events$" "^orders\\.(created|retry|dead)$"add_user 省略密码参数后会交互式提示输入,密码不会进入 shell history 或进程参数。非交互供应任务应让 Secret 管理器把密码直接写入 rabbitmqctl add_user <username> 的标准输入,或者使用受控的 definitions/外部身份后端;不要先把密码拼进命令字符串、CI 参数或日志。RabbitMQ 的访问控制文档列出了交互输入和标准输入两种方式。
权限三元组依次是 configure、write、read,并且是按 vhost 生效的资源名正则。上面的应用账号不能声明或删除资源,只能向 orders.events 写入,并从三个指定队列读取。权限可能按连接或 channel 缓存;变更权限后应让客户端重连再判断是否生效。
拓扑可由受控的部署任务声明。下面使用 management HTTP API 建立 direct exchange、quorum 主队列、classic 延迟重试队列和死信队列。URL 中的 vhost 必须编码;te_orders 不含斜杠,所以可直接使用:
export RMQ_API='http://127.0.0.1:15672/api'
export RMQ_VHOST='te_orders'
umask 077
read -rsp 'orders_topology password: ' RMQ_PASSWORD; echo
RMQ_CURL_CONFIG="$(mktemp)"
RMQ_PASSWORD=${RMQ_PASSWORD//\\/\\\\}
RMQ_PASSWORD=${RMQ_PASSWORD//\"/\\\"}
printf 'user = "orders_topology:%s"\n' "$RMQ_PASSWORD" > "$RMQ_CURL_CONFIG"
unset RMQ_PASSWORD
trap 'rm -f -- "$RMQ_CURL_CONFIG"' EXIT HUP INT TERM
curl --config "$RMQ_CURL_CONFIG" -fsS -H 'content-type: application/json' \
-X PUT "$RMQ_API/exchanges/$RMQ_VHOST/orders.events" \
-d '{"type":"direct","durable":true,"auto_delete":false,"internal":false,"arguments":{}}'
curl --config "$RMQ_CURL_CONFIG" -fsS -H 'content-type: application/json' \
-X PUT "$RMQ_API/exchanges/$RMQ_VHOST/orders.dlx" \
-d '{"type":"direct","durable":true,"auto_delete":false,"internal":false,"arguments":{}}'
curl --config "$RMQ_CURL_CONFIG" -fsS -H 'content-type: application/json' \
-X PUT "$RMQ_API/queues/$RMQ_VHOST/orders.created" \
-d '{"durable":true,"auto_delete":false,"arguments":{"x-queue-type":"quorum"}}'
curl --config "$RMQ_CURL_CONFIG" -fsS -H 'content-type: application/json' \
-X PUT "$RMQ_API/queues/$RMQ_VHOST/orders.dead" \
-d '{"durable":true,"auto_delete":false,"arguments":{"x-queue-type":"quorum"}}'
curl --config "$RMQ_CURL_CONFIG" -fsS -H 'content-type: application/json' \
-X POST "$RMQ_API/bindings/$RMQ_VHOST/e/orders.events/q/orders.created" \
-d '{"routing_key":"order.created","arguments":{}}'
curl --config "$RMQ_CURL_CONFIG" -fsS -H 'content-type: application/json' \
-X POST "$RMQ_API/bindings/$RMQ_VHOST/e/orders.dlx/q/orders.dead" \
-d '{"routing_key":"order.failed","arguments":{}}'
rm -f -- "$RMQ_CURL_CONFIG"
unset RMQ_CURL_CONFIG
trap - EXIT HUP INT TERM临时配置文件由 umask 077 和 mktemp 限制为当前用户可读,curl 的 argv 中只有文件路径;密码既不出现在命令历史里,也不作为 -u 参数暴露。中断处理和显式删除缺一不可。远程管理接口还必须使用 HTTPS,因为 HTTP Basic 只编码凭证,并不加密传输。
队列类型在声明时由 x-queue-type 决定,不能靠 policy 把已有 classic queue 原地变成 quorum queue。生产环境更适合用 policy 管理 DLX、TTL、长度和 delivery limit,因为 policy 可动态调整;把这些值硬编码成 queue argument 会让应用版本与运维策略冲突。
为 quorum 主队列设置死信和毒消息限制:
docker exec te-rabbitmq rabbitmqctl list_feature_flags name state \
| grep stream_queue
docker exec te-rabbitmq rabbitmqctl set_policy -p te_orders orders-quorum-safety \
'^orders\.created$' \
'{"dead-letter-exchange":"orders.dlx","dead-letter-routing-key":"order.failed","dead-letter-strategy":"at-least-once","delivery-limit":5,"delayed-retry-type":"failed","delayed-retry-min":30000,"delayed-retry-max":120000,"max-length-bytes":104857600,"overflow":"reject-publish"}' \
--apply-to quorum_queues --priority 10
docker exec te-rabbitmq rabbitmqctl list_policies -p te_orders
docker exec te-rabbitmq rabbitmqctl list_queues -p te_orders \
name type durable policy messages_ready messages_unacknowledged consumersmax-length-bytes 的 100 MiB 只是开发实验上限,不是生产默认值。生产值要由消息大小分布、到达速率、最大恢复时间和磁盘预算推导。
关键 quorum queue 若要求死信转移过程不丢消息,policy 必须显式设置 dead-letter-strategy=at-least-once、overflow=reject-publish 和 dead-letter-exchange,并确保 stream_queue feature flag 已启用;缺少这些条件会使用或退回默认的 at-most-once 语义。若 feature flag 未启用,先按集群升级流程确认所有节点兼容,再执行 rabbitmqctl enable_feature_flag stream_queue,不要在共享集群看到 disabled 就直接开启。
at-least-once 由源 quorum queue 的内部消费者等待目标 publisher confirm 后再删除死信,因此会增加内存和 CPU 开销,尚未确认的死信仍计入源队列长度限制。DLX 不存在、消息不可路由,或者任一目标队列不可用、达到长度上限而拒绝发布时,源队列会保留消息并周期重试;持续阻塞可能让源队列触顶并拒绝新的发布,恢复时还可能因确认窗口产生重复。生产监控必须同时覆盖源队列、所有死信目标、不可路由事件和 rejected publish。完整限制见 RabbitMQ 4.3 的Quorum Queue 指南。
正向实验:证明消息跨过三次责任转移
用 Python 客户端可以同时验证 mandatory、publisher confirm 和手动 ack。创建隔离环境:
python -m venv .venv-rabbitmq
source .venv-rabbitmq/bin/activate
python -m pip install "pika>=1.3,<2"Windows PowerShell 的激活命令是:
.\.venv-rabbitmq\Scripts\Activate.ps1
python -m pip install "pika>=1.3,<2"创建 publish_once.py:
import json
import os
import uuid
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=os.getenv("RABBITMQ_VHOST", "te_orders"),
credentials=credentials,
heartbeat=30,
blocked_connection_timeout=10,
)
event_id = str(uuid.uuid4())
body = json.dumps({"eventId": event_id, "orderId": "demo-1001"}).encode()
with pika.BlockingConnection(params) as connection:
channel = connection.channel()
channel.confirm_delivery()
confirmed = channel.basic_publish(
exchange="orders.events",
routing_key="order.created",
body=body,
mandatory=True,
properties=pika.BasicProperties(
content_type="application/json",
delivery_mode=pika.DeliveryMode.Persistent,
message_id=event_id,
type="order.created.v1",
),
)
print({"eventId": event_id, "brokerConfirmed": confirmed})设置应用凭证并运行:
export RABBITMQ_USERNAME='orders_app'
export RABBITMQ_PASSWORD='replace-app-password'
export RABBITMQ_VHOST='te_orders'
python publish_once.py预期输出中的 brokerConfirmed 为 True,随后队列的 messages_ready 增加:
docker exec te-rabbitmq rabbitmqctl list_queues -p te_orders \
name messages_ready messages_unacknowledged consumers创建 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=os.getenv("RABBITMQ_VHOST", "te_orders"),
credentials=credentials,
heartbeat=30,
)
with pika.BlockingConnection(params) as connection:
channel = connection.channel()
channel.basic_qos(prefetch_count=10)
method, properties, body = channel.basic_get(
queue="orders.created", auto_ack=False
)
if method is None:
raise SystemExit("queue is empty")
event = json.loads(body)
print({
"eventId": event["eventId"],
"deliveryTag": method.delivery_tag,
"redelivered": method.redelivered,
})
# 真实项目在这里提交本地业务事务或幂等记录。
channel.basic_ack(delivery_tag=method.delivery_tag)运行后再次检查,主队列的 ready 与 unacked 都应回到基线。只有这时,才证明消息完成了 broker 接管、路由入队和消费者确认三段链路。业务正确性还要检查数据库中的幂等记录或领域状态。
反向实验:让错误路由稳定暴露
把 routing key 改成没有 binding 的 order.unknown,保留 mandatory=True:
channel.basic_publish(
exchange="orders.events",
routing_key="order.unknown",
body=b'{"eventId":"negative-routing"}',
mandatory=True,
)Pika 会抛出 pika.exceptions.UnroutableError。这是期望证据:exchange 存在、publish 权限有效,但路由结果为空。若关闭 mandatory,broker 可以确认已经处理 publish,消息却因不可路由被丢弃;这正是“confirm 成功但业务没收到”的典型误判。
排查顺序应固定为:
docker exec te-rabbitmq rabbitmqctl list_exchanges -p te_orders name type durable
docker exec te-rabbitmq rabbitmqctl list_bindings -p te_orders \
source_name source_kind destination_name destination_kind routing_key
docker exec te-rabbitmq rabbitmqctl list_queues -p te_orders \
name messages_ready messages_unacknowledged consumers恢复正确 routing key 后重新发布,不要通过临时创建“拼错名字的队列”掩盖错误。实验产生的错误消息没有进入队列,因此无需 purge。
反向实验:观察 unacked、重投与死信
消费者使用手动 ack 时,delivery 已经从 ready 移到 unacked,但还没有从 broker 删除。把 consume_once.py 的 ack 前加入等待,在另一个终端观察:
input("message is unacked; press Enter to acknowledge")
channel.basic_ack(delivery_tag=method.delivery_tag)此时 messages_unacknowledged=1。直接终止消费者连接,RabbitMQ 会把未确认 delivery 自动重新入队;下一次消费时 redelivered 应为 True。因此消费者必须用 message_id 或业务唯一键做幂等,不能把“RabbitMQ 不重复投递”当作正确性前提。
要验证死信,把失败处理改成拒绝且不重新入队:
channel.basic_nack(delivery_tag=method.delivery_tag, requeue=False)随后检查:
docker exec te-rabbitmq rabbitmqctl list_queues -p te_orders \
name messages_ready messages_unacknowledged consumers预期 orders.created 减少、orders.dead 增加。死信消息会携带 x-death 等头部,消费者可据此判断原队列、原因与次数。当前 policy 使用 at-least-once 转移:死信 exchange 不存在、routing key 无匹配或目标队列不可用时,源队列会保留并周期重试,而不是立即确认删除;恢复后仍可能因确认窗口出现重复。关键死信链应有独立监控、人工处置状态和重放审计。
不要对瞬时错误无限 requeue=True。快速失败会形成重投热循环,消耗 CPU 和网络并饿死正常消息。RabbitMQ 4.3 为 quorum queue 增加了队列内 delayed retry:delayed-retry-min 是必需的最小延迟毫秒数,delayed-retry-max 是可选上限且默认等于最小值,实际延迟为 min(delayed-retry-min × delivery-count, delayed-retry-max)。上面的 policy 从 30s 起步,线性增长并封顶 120s,消息在原 quorum queue 内等待,不需要在 TTL 队列之间反复改写。
delayed-retry-type 决定哪些返回进入延迟:failed 只处理会增加 delivery-count 的失败,例如 AMQP 0.9.1 basic.reject、连接或 channel 携带未确认消息终止;returned 只处理不增加该计数的返回,例如 basic.nack;all 覆盖两类,disabled 为默认值。因此当前 failed policy 适合消费者把可重试失败明确标记为 reject 的约定,不能把 basic.nack(requeue=True) 误认为会按失败次数退避。实体级限流、单行数据库锁等只影响少数消息的瞬时故障适合 delayed retry;若整个数据库离线、所有消息都会失败,应暂停消费者并在依赖恢复后继续,而不是延迟整条积压。跨队列隔离、日历调度和长时间业务编排仍应使用专门的重试拓扑或调度器。无论哪条路径,都要保留业务幂等键、失败分类、最大尝试次数和最终死信状态。参数语义见Quorum Queue 指南和 RabbitMQ 4.3 发布说明。
Publisher confirm、持久化与“至少一次”
生产者调用 basic_publish 返回,不代表 broker 已接管消息。开启 confirm 后,broker 对每条 publish 返回 ack 或 nack。连接在 confirm 到达前断开时,生产者无法确定消息是否已经接管,只能重发未确认消息,因此下游必须容忍重复。
可靠发布还需要同时满足:
exchange 和目标 queue 是 durable;消息使用 persistent delivery mode;生产者开启 confirm 并保存未确认窗口;
不可路由消息使用 mandatory return 或 alternate exchange;confirm 超时后按不确定结果重试,而不是记录“发送成功”;业务状态与待发布事件之间使用 transactional outbox 或等价原子持久化,避免数据库提交成功而 publish 未发生。
RabbitMQ 4.3 的确认机制说明指出,持久消息路由到 durable queue 后通常在持久化时确认;对 quorum queue,则在多数副本接管后确认。Confirm 不能覆盖生产者本地数据库与 publish 之间的原子性,也不能覆盖消费者业务事务。
批量 confirm 比每条同步等待有更高吞吐,但扩大了不确定窗口。生产者应给每条未确认消息保存 sequence、业务键、首次发送时间和重试次数;连接恢复后只重发未确认集合。窗口不能无限增长,应设置最大在途数量和发布超时,通过背压保护应用内存。
Consumer ack 与 prefetch 是一组容量控制
自动 ack 在 broker 写入客户端 socket 后就删除消息。消费者进程随后崩溃时无法恢复,适合可丢弃通知,不适合订单、账务或任务。手动 ack 应发生在业务副作用完成并可恢复之后;如果先 ack 再提交数据库,进程在两者之间崩溃就会丢消息。
Delivery tag 只在当前 channel 内有效,必须在接收它的同一 channel 上 ack。跨线程共享 channel 或重连后使用旧 tag 会触发 unknown delivery tag 并关闭 channel。连接与 channel 恢复后,应用应重新声明拓扑、重新订阅并把未确认 delivery 当作可能重投。
prefetch_count 限制每个消费者允许的未确认消息数。值过小会让消费者等待网络往返,吞吐不足;值过大则把大量消息锁在慢消费者内存里,造成负载不均和故障重投峰值。估算起点可以使用:
prefetch ≈ 单消费者目标吞吐 × 单条处理时间(秒)例如单消费者目标为每秒 40 条,处理 P95 为 0.25 秒,可从 10 开始压测。这个数字只是起点,最终应同时观察 consumer capacity、处理延迟、unacked、进程内存和重投峰值。官方消费者指南提供 consumer capacity 指标,用来判断队列是否经常因为消费者或 prefetch 没有投递空间。
顺序要求会改变并发模型。一个 queue 加一个 active consumer 最容易保持投递顺序,但吞吐和可用性受单消费者限制;多个消费者会并发处理,完成顺序不再稳定。Quorum queue 可使用 Single Active Consumer 在故障时切换 active consumer,但重投和业务重试仍可能改变观察顺序。真正需要按业务键有序时,应把同一键路由到同一串行分片,并让业务状态机拒绝过期版本。
Quorum Queue 的多数派边界
Quorum queue 使用 Raft 复制队列日志,适合长生命周期、数据安全优先的关键队列。官方Quorum Queue 指南建议把三个成员作为实际最小值;增加到五个成员可以容忍两个成员永久不可用,但会增加网络、磁盘和确认延迟。偶数成员不会增加可容忍故障数。
三成员 quorum queue 在一个成员宕机时仍有多数派,可以选主和继续确认;同时失去两个成员时不能形成多数派,队列优先保持一致性而停止推进。此时“其余 RabbitMQ 节点还活着”不等于该队列可用,诊断必须查看每个 quorum queue 的成员和 leader:
docker exec te-rabbitmq rabbitmq-queues quorum_status --vhost te_orders orders.created
docker exec te-rabbitmq rabbitmq-diagnostics cluster_statusQuorum queue 把消息持续写盘并复制,多成员数量和大消息都会放大 I/O。它不适合瞬时 exclusive queue、高频创建删除、极低延迟通知和超长积压;积压达到数百万且需要重复读取时,应评估 RabbitMQ Stream 或 Kafka,而不是继续增加 prefetch。
Classic queue 可以作为非关键临时队列和重试延迟队列,但单节点上的 classic queue 不提供内容复制。旧式 classic mirrored queue 已不再是 RabbitMQ 4.x 的生产方案,迁移时应新建 quorum queue、双写或停机搬迁并核对消息数,不能通过修改参数原地转换。
容量不是只看 messages_ready
队列积压的恢复时间由净消费能力决定:
净排空速率 = 总消费速率 - 到达速率
预计恢复时间 = backlog / 净排空速率若生产每秒 800 条,消费者每秒处理 900 条,十万条积压至少需要约一千秒才能排空;任何处理抖动都会继续延长。若消费速率小于到达速率,增加磁盘只是在推迟故障,不会恢复稳态。
容量评审至少要同时记录消息体 P50/P95/P99、峰值到达率、每个队列消费者数、单条处理 P95、最大可接受积压时间、重试放大倍数、quorum 副本数和磁盘高水位。粗略磁盘预算应包含消息体、索引/日志开销、复制倍数和安全余量,不能只用消息条数乘平均大小。
观察命令:
docker exec te-rabbitmq rabbitmqctl list_queues -p te_orders \
name type messages_ready messages_unacknowledged message_bytes consumers consumer_capacity
docker exec te-rabbitmq rabbitmq-diagnostics alarms
docker exec te-rabbitmq rabbitmq-diagnostics memory_breakdown
docker exec te-rabbitmq rabbitmq-diagnostics disk_spacemessages_ready 是等待投递,messages_unacknowledged 是已经交给消费者但尚未确认。Ready 单调增长通常是消费能力不足、消费者离线或路由突增;unacked 高位不降通常是处理变慢、ack 丢失、prefetch 过大或线程阻塞。内存/磁盘 alarm 会触发发布端流控,生产者表现为延迟升高或连接阻塞,此时应先保护磁盘和消费链,而不是盲目提高 watermark。
项目接入要把不确定结果显式化
应用配置至少区分连接、发布、消费和拓扑:
messaging:
rabbitmq:
addresses: ${RABBITMQ_ADDRESSES:127.0.0.1:5672}
virtual-host: ${RABBITMQ_VHOST:te_orders}
username: ${RABBITMQ_USERNAME}
password: ${RABBITMQ_PASSWORD}
connection-timeout: 5s
requested-heartbeat: 30s
topology:
exchange: orders.events
queue: orders.created
routing-key: order.created
queue-type: quorum
publisher:
confirms: correlated
mandatory: true
max-in-flight: 1000
confirm-timeout: 5s
consumer:
manual-ack: true
prefetch: 10
concurrency: 4
max-processing-time: 30s连接恢复不等于业务恢复。客户端自动重连后,应确认 channel 已重建、publisher confirm 模式已重新开启、consumer 已重新注册、拓扑声明没有属性冲突。若同名 queue 已经以不同的 durable、exclusive 或 x-queue-type 参数存在,声明会关闭 channel;这类错误应该阻止应用进入 ready,而不是降级成日志警告。
消费者处理流程可以固定为:解析并验证 schema,检查幂等键,开启本地事务,执行业务变更并写处理记录,提交事务,最后 ack。可重试错误进入有上限的退避路径;永久格式错误直接 nack 到死信;未知错误保留证据并触发熔断,避免热循环。
故障证据如何反推根因
| 现象 | 第一证据 | 可能根因 | 首个安全动作 |
|---|---|---|---|
| UI 可开、AMQP 失败 | 应用异常与 broker 认证日志 | 把 15672 当 5672、vhost 未编码、无权限 | 用应用身份执行最小连接,不改全局权限 |
| confirm 超时 | pending confirm 数、broker alarm | 网络抖动、磁盘慢、quorum 失去多数派 | 停止扩大在途窗口,保留未确认集合 |
| ready 持续增长 | publish/deliver/ack rate | 消费者离线、处理能力不足、下游慢 | 限流生产或扩消费,计算净排空速率 |
| unacked 持续增长 | channel、consumer、处理耗时 | prefetch 过大、线程阻塞、ack 未执行 | 缩小并发窗口,采集消费者线程证据 |
| 消息重复 | redelivered、x-death、业务幂等记录 | ack 前断连、confirm 未知后重发、重试 | 以业务键去重,不通过 purge 消除现象 |
| 发布成功但队列为零 | return handler、binding 列表 | routing key 无匹配且未处理 mandatory | 修正路由并补 return 监控 |
| 改密码后仍旧行为 | volume、用户列表、节点启动日志 | 旧 volume 已初始化安全状态 | 显式变更用户;仅个人环境才重建 volume |
默认 guest/guest 只允许 loopback 连接。共享环境应创建独立账号,而不是通过 loopback_users = none 放开默认口令。Management UI 的 user tag 只控制 UI 能力,不代替 vhost 的 configure/write/read 权限。
凭证、敏感数据与团队责任
AMQP URI、密码、Erlang cookie、TLS 私钥、OAuth client secret 和内部 broker 地址都不能进入仓库、截图或工单。生产身份至少拆分为拓扑变更、应用发布、应用消费、只读观测和平台管理;凭证轮换要支持双凭证过渡,先创建新身份并验证连接,再滚动应用,最后撤销旧身份。
消息体也可能是敏感数据。管理台的 Get Message、HTTP API 取消息和死信重放都能看到 payload,应限制给受控角色并记录审计。日志只记录 message id、event type、routing key 和追踪标识,不默认输出完整 payload、认证头或个人信息。
团队应为每个 vhost 和 queue 登记 owner、业务等级、生产者、消费者、峰值速率、消息大小、TTL/DLX、幂等键、SLO、告警和清理策略。没有 owner 的队列不得无限保留;临时队列应通过命名与 TTL 自动回收;purge、delete、修改 policy 和重放死信都属于高风险变更。
purge 比 delete 更容易造成错觉:队列和 binding 仍在,消息证据却被清空。共享环境应默认禁止应用账号 purge;确需清理时先停止生产、记录 ready/unacked 数、导出必要证据、确认 owner 和影响面,再只清理指定队列。生产死信重放还应限速、保序评估并保留重放批次号。
清理开发实验而不误伤共享环境
先停止实验消费者,确认没有 unacked,再删除本次创建的资源:
umask 077
read -rsp 'orders_topology password: ' RMQ_PASSWORD; echo
RMQ_CURL_CONFIG="$(mktemp)"
RMQ_PASSWORD=${RMQ_PASSWORD//\\/\\\\}
RMQ_PASSWORD=${RMQ_PASSWORD//\"/\\\"}
printf 'user = "orders_topology:%s"\n' "$RMQ_PASSWORD" > "$RMQ_CURL_CONFIG"
unset RMQ_PASSWORD
trap 'rm -f -- "$RMQ_CURL_CONFIG"' EXIT HUP INT TERM
curl --config "$RMQ_CURL_CONFIG" -fsS -X DELETE \
"$RMQ_API/queues/$RMQ_VHOST/orders.dead"
curl --config "$RMQ_CURL_CONFIG" -fsS -X DELETE \
"$RMQ_API/queues/$RMQ_VHOST/orders.created"
curl --config "$RMQ_CURL_CONFIG" -fsS -X DELETE \
"$RMQ_API/exchanges/$RMQ_VHOST/orders.dlx"
curl --config "$RMQ_CURL_CONFIG" -fsS -X DELETE \
"$RMQ_API/exchanges/$RMQ_VHOST/orders.events"
rm -f -- "$RMQ_CURL_CONFIG"
unset RMQ_CURL_CONFIG
trap - EXIT HUP INT TERM
docker exec te-rabbitmq rabbitmqctl clear_policy -p te_orders orders-quorum-safety
docker exec te-rabbitmq rabbitmqctl delete_user orders_app
docker exec te-rabbitmq rabbitmqctl delete_user orders_topology
docker exec te-rabbitmq rabbitmqctl delete_vhost te_orders个人临时环境可停止容器但保留数据:
docker compose down只有确认 volume 不含其他项目或排障证据后,才执行不可逆清理:
docker compose down -v
rm -rf .venv-rabbitmq共享实例不得执行 down -v,也不得用管理员账号运行批量 delete。清理结果以目标 vhost、用户和资源已不存在为证据,而不是以容器退出为证据。
架构决策的停止条件
RabbitMQ 适合需要灵活路由、低延迟任务分发、每条消息显式确认和有限积压的系统。选择 classic queue 还是 quorum queue,取决于消息是否可重建、是否需要多数派复制以及能否承担额外磁盘和确认延迟;选择单队列多消费者还是按业务键分片,取决于顺序、并发和热点;选择自建还是托管,取决于团队是否能长期承担升级、磁盘、网络分区和故障演练。
当需求变成超长保留、按 offset 反复回放、海量分区顺序流或跨团队事件日志时,应评估 Kafka、Pulsar 或 RabbitMQ Stream。当需求是跨地域 broker 联通,应使用 Federation/Shovel 或应用级复制,不应把 RabbitMQ 集群横跨 WAN。当需求要求数据库变更与消息发送原子一致时,应引入 outbox 与幂等消费,而不是继续堆叠 broker 参数。
最终的上线判断不是“管理台能看到绿色节点”,而是:生产者能区分确认、退回和未知结果;路由拓扑可审查;消费者在崩溃后安全重投;死信可观测、可限速重放;quorum 失去多数派时有明确降级;容量趋势能预测耗尽时间;权限和清理操作不会跨越项目边界。做到这些,RabbitMQ 才从一个能启动的容器变成可治理的消息基础设施。
