上一篇我们讲了 RabbitMQ 的 6 种路由模式,属于"能用";这篇我们讲高级特性,属于"生产级可靠能用"。
🎯 目标:做到 消息 0 丢失 + 业务 0 乱序 + 故障可恢复。
一、全链路消息可靠性(三件套)
要保证消息从生产者发出 → Broker 存储 → 消费者处理 三段不丢失,必须同时启用:
- Producer Confirm:Broker 收到并持久化后 ACK 生产者
- 持久化(Durable):Exchange、Queue、Message 三者都设置为持久化
- Consumer ACK(手动):消费者处理完成再返回 ACK,中途崩了自动重入队

二、Producer Confirm:让生产者放心
默认情况下 RabbitMQ 是"发后即忘",消息到底有没有落到 Broker 磁盘上,生产者不知道。
启用 Confirm 模式(C#)
var factory = new ConnectionFactory { HostName = "localhost" };
using var conn = factory.CreateConnection();
using var ch = conn.CreateModel();
// ★ 开启发布者确认模式(只能在 channel 未被使用时设置一次)
ch.ConfirmSelect();
// 发布消息
byte[] body = Encoding.UTF8.GetBytes("Important data");
ch.BasicPublish(
exchange: "order_events",
routingKey: "created",
mandatory: true, // ★ 找不到队列时走 Return 回调,不静默丢失
basicProperties: props,
body: body
);
// 方式 1:同步等待(简单但吞吐量低)
if (ch.WaitForConfirms(TimeSpan.FromSeconds(5)))
Console.WriteLine("✅ 消息已落到 Broker");
else
throw new Exception("❌ 消息未确认,需要重发");
// 方式 2:异步回调(高吞吐)
ch.BasicAcks += (sender, ea) =>
{
Console.WriteLine($"✅ Broker ACK deliveryTag={ea.DeliveryTag}, multiple={ea.Multiple}");
};
ch.BasicNacks += (sender, ea) =>
{
Console.WriteLine($"❌ Broker NACK deliveryTag={ea.DeliveryTag},需要重发!");
};
ch.BasicReturn += (sender, ea) =>
{
// mandatory=true 触发:找不到可路由的队列时,消息会被退回
Console.WriteLine($"⚠️ Return 消息 {ea.ReplyCode}/{ea.ReplyText}");
};
⚠️ 坑点:mandatory=false 时找不到队列消息会直接被丢弃,生产者完全察觉不到,生产环境务必设为 true 并监听 BasicReturn。
三、持久化三要素
光开了 Confirm 还不够,如果 Broker 宕机重启后队列和消息都没了,等于白搭。必须同时持久化:
| 组件 | 代码 | 含义 |
|---|---|---|
| 交换机持久化 | ExchangeDeclare(..., durable: true) |
重启后交换机还在 |
| 队列持久化 | QueueDeclare(..., durable: true) |
重启后队列定义还在 |
| 消息持久化 | props.Persistent = true |
消息会刷到磁盘,重启不丢 |
三缺一不可。常见错误:队列是 durable,但消息忘了设 Persistent = true,重启后队列空空如也。
四、Consumer ACK:让消费者不丢消息
四种确认方式对比
| 方式 | 代码 | 推荐场景 |
|---|---|---|
| 自动确认(最危险) | BasicConsume(autoAck: true) |
丢消息无所谓的场景,测试用 |
| 手动正向确认 | BasicAck(tag, multiple:false) |
正常消费完成 |
| 手动否定 + 重入队 | BasicNack(tag, multiple:false, requeue:true) |
临时故障,稍后重试 |
| 手动否定 + 丢弃(进死信) | BasicNack(tag, multiple:false, requeue:false) |
消息本身有问题,不要反复重试 |
正确的消费者代码模板
// 1. 一次只取一条(保证公平分发 + 异常时不吞多条)
ch.BasicQos(0, 1, false);
var consumer = new EventingBasicConsumer(ch);
consumer.Received += (model, ea) =>
{
try
{
var msg = Encoding.UTF8.GetString(ea.Body.ToArray());
ProcessBusiness(msg); // 你的业务
// ✅ 处理成功才 ACK
ch.BasicAck(deliveryTag: ea.DeliveryTag, multiple: false);
}
catch (BusinessException bex) // 业务异常(参数不对、格式错)→ 进死信
{
Log.Error(bex, "消息格式错误,丢弃");
ch.BasicNack(deliveryTag: ea.DeliveryTag, multiple: false, requeue: false);
}
catch (DbTimeoutException tex) // 临时异常 → 稍后重试
{
Log.Warn(tex, "DB 超时,重新入队");
ch.BasicNack(deliveryTag: ea.DeliveryTag, multiple: false, requeue: true);
}
catch (Exception ex)
{
Log.Fatal(ex, "未知异常");
// 一定要 ACK / NACK 其中一个,否则 channel 会挂起
ch.BasicNack(deliveryTag: ea.DeliveryTag, multiple: false, requeue: true);
}
};
ch.BasicConsume(queue: "biz_queue", autoAck: false, consumer: consumer); // autoAck 必须 false!
💡 经验:任何时候都别用 autoAck=true 做生产业务,除非你能接受"消费者一崩整批消息全丢"。
五、死信队列(DLX):处理"烂尾消息"

什么是死信?
一条消息进入业务队列后,如果触发以下三种情况之一,就会被转发到 DLX(Dead Letter Exchange),然后流入死信队列:
- 消息被消费者 NACK 且 requeue=false(图路径①)
- 消息的 TTL(过期时间)到期 仍未被消费(图路径②)
- 队列 达到最大长度(max-length),新消息挤掉了头部老消息(图路径③)
死信的价值
- 正常消费者只处理有效消息,不用被"烂数据"阻塞
- 死信队列独立消费,可以做:告警、人工复核、定时重试、写入数据库备查
DLX 配置代码
// 1. 先声明死信交换机 + 死信队列
ch.ExchangeDeclare("dlx_exchange", ExchangeType.Direct, durable: true);
ch.QueueDeclare("dead_letter_queue", durable: true, exclusive: false, autoDelete: false);
ch.QueueBind("dead_letter_queue", "dlx_exchange", routingKey: "");
// 2. 业务队列声明时,在 arguments 中指定 "x-dead-letter-exchange"
var bizArgs = new Dictionary<string, object>
{
{ "x-dead-letter-exchange", "dlx_exchange" },
// 可选:{ "x-dead-letter-routing-key", "some-key" }
// 可选:{ "x-message-ttl", 10 * 60 * 1000 } // 业务队列默认 TTL 10 分钟
// 可选:{ "x-max-length", 10000 } // 队列最大长度
};
ch.QueueDeclare(
queue: "business_queue",
durable: true,
exclusive: false,
autoDelete: false,
arguments: bizArgs // ★ 核心
);
单条消息设置 TTL
var props = ch.CreateBasicProperties();
props.Persistent = true;
props.Expiration = "30000"; // 30 秒 TTL(单位毫秒,字符串!)
ch.BasicPublish(exchange: "", routingKey: "business_queue", basicProperties: props, body: body);
⚠️ 注意:队列级 TTL 是所有消息默认 TTL,消息级 TTL 覆盖队列级。哪个先到按哪个。
六、延迟队列 & 优先级队列

6.1 延迟队列(定时任务神器)
需求:下单 30 分钟未支付 → 自动取消订单。
实现方案有两种:
| 方案 | 实现原理 | 优缺点 |
|---|---|---|
| TTL + DLX(原生) | 消息设 30 分钟 TTL → 业务队列到点过期 → 转发到死信队列 → 死信消费者执行取消 | 无需插件,稳定;但死信路由时间依赖头部扫描 |
| delayed-message-exchange 插件⭐ | 给消息加 header x-delay=毫秒,延迟 Exchange 到期后再投递 |
精度高、语义清晰;但需要装插件 |
推荐使用 延迟插件(官方推荐):
// 安装:rabbitmq-plugins enable rabbitmq_delayed_message_exchange
// 声明类型为 x-delayed-message 的 Exchange
var delayArgs = new Dictionary<string, object> { { "x-delayed-type", "direct" } };
ch.ExchangeDeclare("delay_exchange", "x-delayed-message", durable: true, autoDelete: false, arguments: delayArgs);
ch.QueueDeclare("order_timeout_queue", durable: true);
ch.QueueBind("order_timeout_queue", "delay_exchange", "order.timeout");
// 发布 30 分钟后投递
var props = ch.CreateBasicProperties();
props.Headers = new Dictionary<string, object> { { "x-delay", 30 * 60 * 1000 } }; // 毫秒
ch.BasicPublish("delay_exchange", "order.timeout", props,
body: Encoding.UTF8.GetBytes($"OrderId=1001,请取消订单"));
6.2 优先级队列(插队神器)
需求:VIP 订单优先处理、重要告警优先推送。
// 声明时指定队列最大优先级(0~255,官方建议 1~10 够用)
var args = new Dictionary<string, object> { { "x-max-priority", 10 } };
ch.QueueDeclare("priority_queue", durable: true, exclusive: false, autoDelete: false, arguments: args);
// 发消息时给优先级
var props1 = ch.CreateBasicProperties();
props1.Priority = 10; // P10:最高
ch.BasicPublish(exchange: "", routingKey: "priority_queue", basicProperties: props1, body: Encoding.UTF8.GetBytes("VIP 订单"));
var props2 = ch.CreateBasicProperties();
props2.Priority = 0; // P0:普通
ch.BasicPublish(exchange: "", routingKey: "priority_queue", basicProperties: props2, body: Encoding.UTF8.GetBytes("普通订单"));
// ★ 消费者侧:高优先级消息会被先推出来
// 注意:消费者空闲时才会从队列取高优,若消费者一直忙碌则按 prefetch 批量拉的优先级排序
⚠️ 优先级队列只对队列中排队的消息有效,如果消费者都在忙且 prefetch 拉走了一批低优消息,这批不会被高优插队。解决办法:减小
prefetchCount到 1~2。
七、惰性队列(Lazy Queue,百万级积压必备)
场景:队列里积压了几十万甚至上百万条消息,普通队列会把消息尽量放内存,然后分页到磁盘,内存压力很大。
惰性队列:一上来就把消息写磁盘,消费时才分页加载到内存。适合"积压兜底"的队列。
var args = new Dictionary<string, object> { { "x-queue-mode", "lazy" } };
ch.QueueDeclare("back_pressure_queue", durable: true, arguments: args);
代价:吞吐会稍微低一些(多了磁盘 IO),但内存占用骤降,RabbitMQ 不容易 OOM。生产建议所有可能积压的队列都设 lazy。
八、生产级完整配置模板总结
var args = new Dictionary<string, object>
{
{ "x-dead-letter-exchange", "dlx_exchange" }, // 绑定死信交换机
{ "x-message-ttl", 15 * 60 * 1000 }, // 消息默认 15 分钟 TTL
{ "x-max-length", 500_000 }, // 队列最多 50w 条(防爆)
{ "x-queue-mode", "lazy" }, // 惰性队列,积压不怕
{ "x-max-priority", 10 } // 启用优先级(按需)
};
ch.QueueDeclare("biz_queue_production", durable: true, exclusive: false, autoDelete: false, arguments: args);
配合:
- Producer 端:
ConfirmSelect+mandatory=true+ 重试机制 - Consumer 端:
BasicQos(1)+autoAck=false+ try/catch 分类 ACK/NACK - 运维端:监控队列长度、死信数量、Consumer 连接数
做到这些,你的 RabbitMQ 就能在生产稳定跑起来了。🎉