RabbitMQ 运行链:路由、Confirm、Ack 与 Prefetch 的真实边界
贯通 exchange、binding、queue、publisher confirm、consumer ack、prefetch、重投与死信,解释消息所有权如何逐段转移。
RabbitMQ 官方 Publisher Confirms 明确区分 publisher confirms 与 consumer acknowledgements;Consumer Prefetch 说明预取限制作用于尚未确认的投递。
Exchange 不保存业务成功,只执行路由决定
生产者把消息发布到 exchange,并用 routing key 参与 binding 匹配。direct 按精确键,topic 按分段模式,fanout 忽略 routing key,headers 根据头字段匹配。路由结果可能是零个、一个或多个队列。零路由并不等于网络失败:未启用 mandatory 时,消息可能被正常接收后直接丢弃;启用 mandatory 才能收到 returned message,alternate exchange 则把未路由消息交给另一套拓扑。
拓扑声明需要被当作契约。exchange 类型、durable、auto-delete、arguments 与已有实体不一致时,声明会关闭 channel,而不是替你修改线上资源。队列名称、死信参数和仲裁队列类型的变更应通过新建、迁移、切换完成。把应用启动时的 declare 当作任意 schema migration,会让一次配置漂移演变为全实例启动失败。
Confirm 的证据范围由队列类型和持久化策略决定
publisher confirm 是 channel 上的异步序号协议。Broker 可以批量确认多个 delivery sequence;客户端必须维护在途集合,处理单个或 multiple ack/nack,并限制最大在途数量。同步逐条等待最容易理解,却会把吞吐限制为往返时延的倒数;无限异步发送则会在 Broker 变慢时把消息堆进客户端内存。
confirm 只说明 Broker 完成了该消息在发布侧所要求的处理。它不证明某个消费者看见消息,更不证明业务数据库已提交。连接在消息写出后、confirm 到达前断开时,结果是不确定而不是失败;生产者若重发,Broker 可能收到重复。因此发布端仍需稳定 eventId,消费端仍需幂等。
Ack 决定队列是否仍拥有重投责任
自动 ack 在投递后立即移除队列责任,适合允许丢失且处理极短的观察性数据;业务消息通常使用手动 ack,在本地事务成功后确认。basic.reject 只能处理单条,basic.nack 可批量;requeue=true 会让消息重新竞争,若故障不可恢复就会形成高频重投循环。requeue=false 配合 DLX 才能把毒消息移出主队列,但死信路由本身也要验证容量与可用性。
prefetch 约束一个消费者同时持有的未确认消息数。值过大时单个消费者囤积消息,故障重连造成批量重投;值过小时吞吐被往返与处理时间限制。合理值来自每条消息内存、处理并发、下游连接池和可接受重投批次,而不是固定经验数字。
用所有权转移图读懂可靠性
失败矩阵比“保证不丢”更可执行
工程选择的核心是把不可避免的不确定窗口转成可重复、可查询、可补偿的状态。不要用无限重试掩盖未知结果,也不要用提前确认换取表面低 lag。
Channel 与连接故障为什么会放大重复
连接承载 TCP 与心跳,channel 在其上复用协议状态。大量消费者共用一条连接可降低连接数,却会让连接抖动同时撤销多条 channel 上的未确认投递;每线程一连接又会放大文件描述符、TLS 和心跳成本。发布连接与消费连接通常隔离,关键队列再按故障域拆分,避免慢消费者影响发布 confirm。
客户端自动恢复可以重建拓扑与订阅,但进程内 delivery tag、confirm sequence 和未完成 Future 不能跨 channel 世代复用。恢复逻辑必须把旧世代结果判为失效,重新查询业务状态或依赖 eventId 重投;否则旧 confirm 可能错误完成新请求,旧 ack 也可能作用于不存在的投递标签。
用两个 Java 状态模型固定不变量
下面的模型不连接真实 Broker。它们把本篇最容易混淆的状态、位置或所有权变化压缩成确定输出;随后再用 RabbitMQ、Kafka 或 RocketMQ 的集成环境验证协议、持久化、重平衡和故障时序。
javac --release 17 -Xlint:all -Werror examples/backend-development/message-event/rabbitmq-runtime/RoutingConfirmDemo.java examples/backend-development/message-event/rabbitmq-runtime/ConsumerAckPrefetchDemo.java
java -cp examples/backend-development/message-event/rabbitmq-runtime RoutingConfirmDemo
java -cp examples/backend-development/message-event/rabbitmq-runtime ConsumerAckPrefetchDemoroutedQueues=[billing.q, audit.q] mandatoryReturn=false confirmed=true
prefetch=2 unacked=2 thirdDeliveryBlocked=true afterAckUnacked=1 redelivered=true模型输出变化时,应定位是哪条业务不变量被修改,而不是只更新示例字符串。真实集成测试还要注入进程崩溃、响应丢失、Broker 切换和下游超时,因为这些窗口无法由纯内存模型模拟。
契约、权限与数据生命周期不能交给默认值
用分层测试证明协议与业务同时成立
把设计决定写成可以被证伪的承诺
指标、日志与容量要围绕状态迁移
容量模型至少计算平均/峰值每秒消息数、平均/p99 字节、保留时间、副本倍数、重试放大、批处理密度、消费者处理时长与下游并发。恢复测试要在稳定生产流量仍存在时注入积压,验证净排空速率,而不是暂停生产后测一个理想峰值。
上线与演练以可恢复为验收标准
故障演练依次注入生产响应丢失、Broker 节点切换、消费者提交后崩溃、下游超时、永久坏消息和热点 key。每次演练都回答三件事:业务不变量是否保持、重复或延迟是否有持久证据、系统是否在容量预算内自动或人工收敛。
