上一篇我们讲了 RabbitMQ 的 6 种路由模式,属于"能用";这篇我们讲高级特性,属于"生产级可靠能用"。

🎯 目标:做到 消息 0 丢失 + 业务 0 乱序 + 故障可恢复


一、全链路消息可靠性(三件套)

要保证消息从生产者发出 → Broker 存储 → 消费者处理 三段不丢失,必须同时启用:

  1. Producer Confirm:Broker 收到并持久化后 ACK 生产者
  2. 持久化(Durable):Exchange、Queue、Message 三者都设置为持久化
  3. Consumer ACK(手动):消费者处理完成再返回 ACK,中途崩了自动重入队

RabbitMQ 全链路 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),然后流入死信队列

  1. 消息被消费者 NACK 且 requeue=false(图路径①)
  2. 消息的 TTL(过期时间)到期 仍未被消费(图路径②)
  3. 队列 达到最大长度(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 就能在生产稳定跑起来了。🎉