惯性聚合 高效追踪和阅读你感兴趣的博客、新闻、科技资讯
阅读原文 在惯性聚合中打开

推荐订阅源

V
Visual Studio Blog
I
InfoQ
freeCodeCamp Programming Tutorials: Python, JavaScript, Git & More
博客园 - 【当耐特】
小众软件
小众软件
B
Blog RSS Feed
大猫的无限游戏
大猫的无限游戏
博客园 - 三生石上(FineUI控件)
Engineering at Meta
Engineering at Meta
人人都是产品经理
人人都是产品经理
Microsoft Security Blog
Microsoft Security Blog
Last Week in AI
Last Week in AI
H
Help Net Security
爱范儿
爱范儿
云风的 BLOG
云风的 BLOG
博客园 - 司徒正美
Y
Y Combinator Blog
H
Hackread – Cybersecurity News, Data Breaches, AI and More
Microsoft Azure Blog
Microsoft Azure Blog
L
LangChain Blog
WordPress大学
WordPress大学
GbyAI
GbyAI
Google DeepMind News
Google DeepMind News
腾讯CDC

博客园 - Hey,Coder!

微信小程序 web方案实现调音器功能 systemctl slice配置docker最大占用cpu 内存 pipenv环境变量配置 .net8验证码 数据库差异对比工具 docker部署prometheus grafana pgexporter Alertmanager ubuntu24 换源 ssh配置 ubuntu24 Harbor部署 ubuntu24 kvm部署 cockpit web管理界面 .net api代理 Ubuntu 24.04 DMAR IOMMU 报错修复 uniapp iOS App 上架 vs文件软链接 Docker 网络别名 .net8通用中间件处理公共返回值结构 .net环境下OpenTelemetry Body 读取注意事项 jaeger界面显示异常 .net opentelemetry collector jaeger c# opentelemetry 自定义export c# opentelemetry 自定义export c# 生成对象的初始化代码 Elasticsearch ILM 与 Data Stream 概念及实例说明文档 portainer httphelper封装 Ocelot + Consul + SignalR 粘性会话 正交软件架构 docker postgresql17 主从复制 docker rabbitmq quartz dashboard docker consul registrator 服务自动注册consul 微服务动态扩容
rabbitmq总结与c#示例
Hey,Coder! · 2026-06-15 · via 博客园 - Hey,Coder!

RabbitMQ 核心概念与实践总结

1. 基础架构与消息流转

  • 消息永远只能发送给 Exchange,不能直接发送给 Queue。
  • Exchange 根据绑定规则将消息路由到一个或多个 Queue。
  • Queue 存储消息,供消费者消费。

默认交换机(Default Exchange)

  • 名称为空字符串 "",类型为 direct
  • 每个 Queue 自动绑定到默认交换机,绑定键(routing key)等于队列名。
  • 使用 BasicPublish("", "queue_name", body) 等效于直接发送给队列。

2. Exchange 类型

类型 路由规则 典型用途
direct 精确匹配 routing key 任务分发、点对点
topic 通配符匹配(* 匹配一段,# 匹配多段) 事件驱动、发布订阅
fanout 忽略 routing key,广播给所有绑定的 Queue 广播通知、缓存刷新
headers 按消息头(headers)匹配,x-match=all/any 复杂路由(较少用)

3. Virtual Host(vhost)

  • 逻辑隔离命名空间:不同 vhost 下的 Exchange、Queue、Binding 互不可见。
  • 权限按 user + vhost 配置(configure / write / read)。
  • 典型用法:环境隔离(/prod/staging)、业务隔离(/order-system/payment-system)。
  • 注意:vhost 不隔离 CPU/内存/磁盘,仅为管理隔离。

4. Queue 属性详解

属性 含义
Durable 队列元数据持久化(重启后队列仍在,消息是否持久化看消息属性)
Exclusive 队列仅对创建它的 Connection 可见,Connection 关闭时自动删除
AutoDelete 最后一个消费者取消订阅后自动删除队列
Arguments 扩展参数,如 x-dead-letter-exchangex-message-ttlx-max-length

Exclusive vs AutoDelete

  • exclusive=true:队列绑定到创建它的连接,连接断开即删,其他连接无法使用。
  • autoDelete=true:队列在没有消费者时自动删除(前提是有过消费者)。

5. 死信队列(DLX / Dead Letter Exchange)

触发条件(官方定义)

  1. 消费者使用 basic.rejectbasic.nackrequeue=false
  2. 消息 TTL 过期。
  3. 队列达到最大长度(x-max-length / x-max-length-bytes)。
  4. 消息被超过 delivery-limit 次返回(仅 Quorum Queue)。

配置步骤

  1. 创建死信交换机(普通 Exchange,如 my.dlx,类型建议 fanoutdirect)。
  2. 创建死信队列(如 my.queue.dlq)并绑定到死信交换机。
  3. 在业务队列的 Arguments 中设置:
    • x-dead-letter-exchange = my.dlx
    • (可选)x-dead-letter-routing-key = 自定义 routing key

常见误区

  • Dashboard 的 Get Messages → Reject 不会触发死信(该操作通过 HTTP API 实现,不属于标准 consumer 路径,消息会被直接丢弃)。
  • 死信交换机不是自动绑定的,需要手动将 DLQ 绑定到 DLX。

6. NACK vs Reject

方法 支持批量 参数
BasicReject(deliveryTag, requeue) ❌ 仅单条 deliveryTag, requeue
BasicNack(deliveryTag, multiple, requeue) ✅ 通过 multiple=true 批量 deliveryTag, multiple, requeue
  • 共同点requeue=true 重新入队,requeue=false 触发死信(如果配置了 DLX)或丢弃。
  • 推荐:统一使用 BasicNackmultiple=false 时等价于 BasicReject)。

7. 消费频率控制

第一层:Prefetch(QoS)

  • 控制消费者未确认消息的最大数量,防止内存爆炸。
  • 必须搭配 autoAck=false
  • 语法:channel.BasicQos(0, prefetchCount, global: false)
  • 典型取值:处理慢的任务用 1~3,常规业务 5~20,轻量任务 30~100。

第二层:应用层限速(RabbitMQ 不原生支持)

  • 使用令牌桶、漏桶或简单计数器实现“每秒 N 条”的硬限速。
  • 示例:RateLimiter 类 + TryAcquire() 判断是否继续消费。

8. Producer 是否应该知道 Queue?

理论原则

  • Producer 只应知道 Exchange,不应知道 Queue(解耦)。
  • Queue 由消费者或运维创建并绑定到 Exchange。

工程现实

  • 为确保消息不丢,许多项目让 Producer 也声明 Queue(尤其是没有独立基础设施团队时)。
  • 风险:Producer 和 Consumer 声明参数不一致导致 406 PRECONDITION_FAILED
  • 最佳折衷
    • 使用共享常量/配置文件统一管理 Queue 声明参数。
    • 或通过 Policy 管理死信等参数,代码中只声明必要属性。

9. 防止消息路由丢失:mandatory 与 Alternate Exchange

mandatory 标志

  • 设置 BasicPublishmandatory=true
  • 当消息无法路由到任何 Queue 时,RabbitMQ 通过 basic.return 将消息返还给生产者。
  • 生产者需注册 BasicReturn 事件处理退回消息。
  • 适用:需要实时感知路由失败并自主处理。

Alternate Exchange(备用交换机)

  • 在主 Exchange 的 Arguments 中设置 alternate-exchange
  • 当主 Exchange 路由失败时,消息自动转发到备用 Exchange。
  • 备用 Exchange 通常绑定一个“兜底队列”用于存储未路由消息。
  • 适用:希望自动兜底、无需生产者保持连接。

对比与组合

特性 mandatory Alternate Exchange
处理主体 生产者(回调) RabbitMQ 自动转发
是否需要长连接
典型用途 实时告警、重试 持久化兜底、审计

推荐组合:Alternate Exchange 作为第一道防线,mandatory 作为可选补充。


10. 常见错误与排查

死信队列始终为空

  • 业务队列 Arguments 中是否设置了 x-dead-letter-exchange
  • 死信交换机是否已绑定到死信队列?
  • 消费者是否使用了 requeue=false?(Dashboard 的 Get Messages 不会触发死信)
  • 死信交换机类型与绑定 routing key 是否匹配?

406 PRECONDITION_FAILED

  • 队列已存在,但声明参数(durable、exclusive、autoDelete、arguments)不一致。
  • 解决:删除队列重建,或使用 Policy 统一管理。

消息丢失

  • 未设置 mandatory 且无 Alternate Exchange,且无队列绑定。
  • 消费者使用 autoAck=true 且处理异常未手动 NACK。
  • 队列设置为 autoDelete 且消费者断开。

11. 最佳实践速查表

场景 推荐做法
简单任务队列(点对点) 默认交换机 + 一个 Queue
事件驱动(多个消费者) Topic Exchange + 多个 Queue(每个消费者独立绑定)
广播通知 Fanout Exchange + 多个 Queue
失败重试 / 延迟处理 死信交换机(DLX)+ TTL 或延迟插件
生产环境消息不丢 Alternate Exchange + 持久化 Queue + 手动 ACK
多环境隔离 不同 vhost(/prod, /staging, /dev)
消费者限流 BasicQos(prefetchCount=N) + autoAck=false
精确限速(N条/秒) 应用层令牌桶 + Prefetch 辅助

本文档覆盖了 RabbitMQ 从入门到生产部署的核心知识点,可作为团队内部参考手册。

代码示例

gitee地址

 public enum ExchangeTypeEnum
 {
     Direct,
     Fanout,
     Topic,
     Headers
 }

 public static class ExchangeTypeExtensions
 {
     public static string ToRabbitMQString(this ExchangeTypeEnum type)
     {
         return type switch
         {
             ExchangeTypeEnum.Direct => "direct",
             ExchangeTypeEnum.Fanout => "fanout",
             ExchangeTypeEnum.Topic => "topic",
             ExchangeTypeEnum.Headers => "headers",
             _ => "direct"
         };
     }
 }
public class RabbitMQOptions
{
    /// <summary>
    /// RabbitMQ服务器主机名
    /// </summary>
    public string HostName { get; set; } = "localhost";

    /// <summary>
    /// RabbitMQ服务器端口
    /// </summary>
    public int Port { get; set; } = 5672;

    /// <summary>
    /// 用户名
    /// </summary>
    public string UserName { get; set; } = "guest";

    /// <summary>
    /// 密码
    /// </summary>
    public string Password { get; set; } = "guest";

    /// <summary>
    /// 虚拟主机
    /// </summary>
    public string VirtualHost { get; set; } = "/";

    /// <summary>
    /// 交换机名称
    /// </summary>
    public string ExchangeName { get; set; } = string.Empty;

    /// <summary>
    /// 交换机类型,默认 Direct
    /// </summary>
    public ExchangeTypeEnum ExchangeType { get; set; } = ExchangeTypeEnum.Direct;

    /// <summary>
    /// 队列名称
    /// </summary>
    public string QueueName { get; set; } = string.Empty;

    /// <summary>
    /// 路由键
    /// </summary>
    public string RoutingKey { get; set; } = string.Empty;

    /// <summary>
    /// 是否自动删除
    /// </summary>
    public bool AutoDelete { get; set; } = false;

    /// <summary>
    /// 是否持久化
    /// </summary>
    public bool Durable { get; set; } = true;

    /// <summary>
    /// 是否排他
    /// </summary>
    public bool Exclusive { get; set; } = false;

    /// <summary>
    /// 队列/交换机参数
    /// </summary>
    public IDictionary<string, object?>? Arguments { get; set; }

    /// <summary>
    /// 死信队列名称,为空时自动生成 {QueueName}_dlq
    /// </summary>
    public string DeadLetterQueueName { get; set; } = string.Empty;

    /// <summary>
    /// 预取消息数量,0表示不限制
    /// </summary>
    public ushort PrefetchCount { get; set; } = 0;

    /// <summary>
    /// 预取消息总大小(字节),0表示不限制
    /// </summary>
    public uint PrefetchSize { get; set; } = 0;

    /// <summary>
    /// QoS是否对整个通道生效,false仅对当前消费者生效
    /// </summary>
    public bool GlobalQoS { get; set; } = false;
}

public abstract class RabbitMQClientBase : IAsyncDisposable
{
    protected readonly RabbitMQOptions _options;
    protected IConnection? _connection;
    protected IChannel? _channel;
    private bool _disposed;

    protected RabbitMQClientBase(RabbitMQOptions options)
    {
        _options = options ?? throw new ArgumentNullException(nameof(options));
    }

    protected virtual async ValueTask EnsureConnectionAsync()
    {
        if (_connection == null || !_connection.IsOpen)
        {
            var factory = new ConnectionFactory
            {
                HostName = _options.HostName,
                Port = _options.Port,
                UserName = _options.UserName,
                Password = _options.Password,
                VirtualHost = _options.VirtualHost
            };
            _connection = await factory.CreateConnectionAsync();
        }
    }

    protected virtual async ValueTask EnsureChannelAsync()
    {
        await EnsureConnectionAsync();
        if (_channel != null && !_channel.IsClosed)
        {
            return;
        }
        _channel = await _connection!.CreateChannelAsync();
        await ConfigureChannelAsync();
    }

    protected virtual async ValueTask ConfigureChannelAsync()
    {
        if (string.IsNullOrEmpty(_options.QueueName))
        {
            return;
        }
        await _channel!.QueueDeclareAsync(
            queue: _options.QueueName,
            durable: _options.Durable,
            exclusive: _options.Exclusive,
            autoDelete: _options.AutoDelete,
            arguments: _options.Arguments
        );

        await DeclareDeadLetterExchangeAndQueueAsync();
    }

    private async ValueTask DeclareDeadLetterExchangeAndQueueAsync()
    {
        if (_options.Arguments == null || !_options.Arguments.TryGetValue("x-dead-letter-exchange", out var dlxValue))
        {
            return;
        }
        var dlxName = dlxValue as string;
        if (string.IsNullOrEmpty(dlxName))
        {
            return;
        }
        await _channel!.ExchangeDeclareAsync(
            exchange: dlxName,
            type: ExchangeType.Direct,
            durable: _options.Durable,
            autoDelete: _options.AutoDelete
        );

        var dlqName = !string.IsNullOrEmpty(_options.DeadLetterQueueName) ? _options.DeadLetterQueueName : $"{_options.QueueName}_dlq";
        await _channel.QueueDeclareAsync(
            queue: dlqName,
            durable: _options.Durable,
            exclusive: _options.Exclusive,
            autoDelete: _options.AutoDelete
        );

        string bindingKey;
        if (_options.Arguments.TryGetValue("x-dead-letter-routing-key", out var dlrkValue))
        {
            bindingKey = dlrkValue as string ?? _options.QueueName;
        }
        else
        {
            bindingKey = !string.IsNullOrEmpty(_options.RoutingKey) ? _options.RoutingKey : _options.QueueName;
        }

        await _channel.QueueBindAsync(
            queue: dlqName,
            exchange: dlxName,
            routingKey: bindingKey
        );
    }

    protected virtual async ValueTask DisposeAsync(bool disposing)
    {
        if (_disposed) return;

        if (disposing)
        {
            if (_channel != null)
            {
                await _channel.DisposeAsync();
            }
            if (_connection != null)
            {
                await _connection.DisposeAsync();
            }
        }

        _disposed = true;
    }

    public async ValueTask DisposeAsync()
    {
        await DisposeAsync(true);
        GC.SuppressFinalize(this);
    }
}

  public abstract class RabbitMQConsumerBase : RabbitMQClientBase
  {
      protected AsyncEventingBasicConsumer? _consumer;
      private string? _consumerTag;

      protected RabbitMQConsumerBase(RabbitMQOptions options) : base(options) { }

      public virtual async Task StartConsumingAsync()
      {
          await EnsureChannelAsync();

          if (_options.PrefetchCount > 0 || _options.PrefetchSize > 0)
          {
              await _channel!.BasicQosAsync(_options.PrefetchSize, _options.PrefetchCount, _options.GlobalQoS);
          }

          _consumer = new AsyncEventingBasicConsumer(_channel!);
          _consumer.ReceivedAsync += OnMessageReceivedAsync;

          _consumerTag = await _channel!.BasicConsumeAsync(
              queue: _options.QueueName,
              autoAck: false,
              consumer: _consumer,
              arguments: _options.Arguments,
              consumerTag: ""
          );
      }

      protected virtual async Task OnMessageReceivedAsync(object? sender, BasicDeliverEventArgs e)
      {
          try
          {
              var message = Encoding.UTF8.GetString(e.Body.ToArray());
              await ProcessMessageAsync(message, e.DeliveryTag);
          }
          catch (Exception ex)
          {
              await HandleErrorAsync(ex, e);
          }
      }

      protected abstract Task ProcessMessageAsync(string message, ulong deliveryTag);

      protected virtual async Task HandleErrorAsync(Exception ex, BasicDeliverEventArgs e)
      {
          await _channel!.BasicNackAsync(e.DeliveryTag, false, false);
      }

      protected virtual async Task AckMessageAsync(ulong deliveryTag)
      {
          await _channel!.BasicAckAsync(deliveryTag, false);
      }

      protected virtual async Task NackMessageAsync(ulong deliveryTag, bool requeue = false)
      {
          await _channel!.BasicNackAsync(deliveryTag, false, requeue);
      }

      public virtual async Task StopConsumingAsync()
      {
          if (!string.IsNullOrEmpty(_consumerTag) && _channel != null && _channel.IsOpen)
          {
              await _channel.BasicCancelAsync(_consumerTag);
              _consumerTag = null;
          }
      }

      protected override async ValueTask DisposeAsync(bool disposing)
      {
          if (disposing)
          {
              await StopConsumingAsync();
          }
          await base.DisposeAsync(disposing);
      }
  }

 public abstract class RabbitMQProducerBase : RabbitMQClientBase
 {
     protected RabbitMQProducerBase(RabbitMQOptions options) : base(options) { }

     protected override async ValueTask ConfigureChannelAsync()
     {
         await base.ConfigureChannelAsync();

         if (!string.IsNullOrEmpty(_options.ExchangeName))
         {
             await _channel!.ExchangeDeclareAsync(
                 exchange: _options.ExchangeName,
                 type: _options.ExchangeType.ToRabbitMQString(),
                 durable: _options.Durable,
                 autoDelete: _options.AutoDelete,
                 arguments: _options.Arguments
             );

             if (!string.IsNullOrEmpty(_options.QueueName) && !string.IsNullOrEmpty(_options.RoutingKey))
             {
                 await _channel.QueueBindAsync(
                     queue: _options.QueueName,
                     exchange: _options.ExchangeName,
                     routingKey: _options.RoutingKey,
                     arguments: _options.Arguments
                 );
             }
         }
     }

     public virtual async Task PublishAsync(string message)
     {
         await EnsureChannelAsync();
         var body = Encoding.UTF8.GetBytes(message);
         await PublishMessageAsync(body);
     }

     private async Task PublishMessageAsync(ReadOnlyMemory<byte> body)
     {
         var props = new BasicProperties { Persistent = _options.Durable };

         if (!string.IsNullOrEmpty(_options.ExchangeName))
         {
             await _channel!.BasicPublishAsync(_options.ExchangeName, _options.RoutingKey, false, props, body);
         }
         else if (!string.IsNullOrEmpty(_options.QueueName))
         {
             await _channel!.BasicPublishAsync(string.Empty, _options.QueueName, false, props, body);
         }
     }
 }