在日常 .NET 项目里接入 RabbitMQ,最繁琐的事情不是”发消息”那一行 channel.BasicPublish,而是前置的一长串模板代码:

var factory = new ConnectionFactory { HostName = "localhost" };
using var connection = factory.CreateConnection();
using var channel = connection.CreateModel();

channel.ExchangeDeclare(exchange: "switchdata.direct", type: ExchangeType.Direct);
channel.QueueDeclare(queue: "switchdata.order", durable: true, ...);
channel.QueueBind(queue: "switchdata.order", exchange: "switchdata.direct", routingKey: "order.created");

var body = Encoding.UTF8.GetBytes(json);
channel.BasicPublish(exchange: "switchdata.direct", routingKey: "order.created", ..., body: body);

一个项目里如果有十个消息类型,这段模板就会复制十次。交换机/队列名改一处、死信配置加一项、延迟消息开一关——所有拓扑细节散落在各处,改起来心惊胆战。

今天我来把我们在生产项目(SwitchData)里落地的一套 RabbitMQ 客户端封装完整拆解出来:用特性(Attribute)做消息模型标注,用泛型做类型驱动,用 ConcurrentDictionary 做拓扑缓存,用 SemaphoreSlim 做连接延迟创建,最后把”发送一条消息”的调用压缩成一行泛型方法。

这套封装最终形成四个文件:

文件 职责
RabbitMqConfig.cs 连接配置,host、port、用户、心跳、自动恢复开关
RabbitMqQueueAttribute.cs 标注在消息模型类上,声明交换机/队列/路由键/死信/延迟
RabbitMqClient.cs 核心客户端:单例连接、拓扑缓存、泛型 Send/Receive、资源释放
RabbitMqExchangeType 枚举交换机类型,含 Delayed 映射到 x-delayed-message 插件类型

一、整体架构

先看一眼消息从模型类到 RabbitMQ 的完整链路:

flowchart TD
    A["消息模型类<br/>标记 [RabbitMqQueue] 特性"] --> B["RabbitMqClient.SendMessageAsync<T>"]
    B --> C{_attributeCache<br/>反射缓存}
    C -- 命中 --> D["BuildProperties / BuildExchangeArgs / BuildQueueArgs<br/>组装死信、延迟参数"]
    C -- 未命中 --> E["GetCustomAttribute + Validate<br/>首次提取并缓存"]
    E --> D
    D --> F{"_topologyCache<br/>拓扑声明缓存"}
    F -- 命中 --> G["直接创建 Channel"]
    F -- 未命中 --> H["ExchangeDeclare + QueueDeclare + QueueBind<br/>如有死信则一并声明"]
    H --> G
    G --> I["channel.BasicPublishAsync"]

    subgraph RabbitMqClient 内部
        C
        F
        H
    end

发送端和消费端共享同一个特性缓存和同一个拓扑缓存,拓扑声明幂等——重复调用只会声明一次,不会对 RabbitMQ 造成额外压力。


二、RabbitMqQueueAttribute:把拓扑写在类上

先解决”消息模型长什么样”的问题。我们给每个消息 DTO 打一个 [RabbitMqQueue] 特性,把交换机、队列、路由键这些声明式配置全部塞进去:

[RabbitMqQueue(
    ExchangeName = "switchdata.direct",
    ExchangeType = RabbitMqExchangeType.Direct,
    QueueName = "switchdata.order",
    RoutingKey = "order.created",
    Durable = true,
    // 死信队列
    DeadExchangeName = "switchdata.dead.direct",
    DeadQueueName = "switchdata.dead.order",
    DeadRoutingKey = "order.dead")]
public class OrderCreatedEvent
{
    public long OrderId { get; set; }
    public string CustomerId { get; set; }
    public decimal Amount { get; set; }
    public DateTime CreatedAt { get; set; }
}

延迟队列消息只需要改一行:把 ExchangeType = RabbitMqExchangeType.Delayed,然后调用发送方法时传 delayMilliseconds 参数即可。

特性类的实现有几个值得看的细节:

1. Delayed 类型自动映射到 x-delayed-message

RabbitMQ 的延迟消息需要安装 rabbitmq_delayed_message_exchange 插件,插件要求交换机类型字符串必须是 x-delayed-message,而不是我们在 .NET 里用枚举写的 Delayed。我们在枚举和字符串之间做了一层自动映射:

public static string GetExchangeTypeName(RabbitMqExchangeType exchangeType)
{
    if (exchangeType == RabbitMqExchangeType.Delayed)
        return "x-delayed-message";   // 插件要求的硬编码类型字符串
    return exchangeType.ToString().ToLower();
}

2. 属性 setter 里的自动赋值

ExchangeType 的 setter 不是单纯赋值,而是同步更新 ExchangeTypeName:

public RabbitMqExchangeType ExchangeType
{
    get => _exchangeType;
    set
    {
        if (_exchangeType != value)
        {
            _exchangeType = value;
            ExchangeTypeName = GetExchangeTypeName(_exchangeType);
        }
    }
}

同样的模式也出现在 DelayedExchangeType 和 DeadExchangeType 上。这样外部拿到的 ExchangeTypeName 永远是经过映射后的正确字符串,使用者不需要关心底层插件细节。

3. Validate:构造时的契约检查

特性自身就带着校验逻辑,RabbitMqClient 取到特性之后立刻调 .Validate(type):

public void Validate(Type type)
{
    if (string.IsNullOrWhiteSpace(ExchangeName))
        throw new InvalidOperationException($"{type.Name}:ExchangeName 不能为空");
    // DeadExchangeName 不为空时,DeadQueueName 和 DeadRoutingKey 也必须齐全
    if (!string.IsNullOrWhiteSpace(DeadExchangeName))
    {
        if (string.IsNullOrWhiteSpace(DeadQueueName))
            throw new InvalidOperationException($"{type.Name}:DeadQueueName 不能为空");
        if (string.IsNullOrWhiteSpace(DeadRoutingKey))
            throw new InvalidOperationException($"{type.Name}:DeadRoutingKey 不能为空");
    }
}

契约式编程的好处是:消息模型类在首次发送前就把好关,而不是等到发消息时才报一个模糊的 “queue not declared”。


三、RabbitMqClient:核心客户端的七个关键设计

设计一:单例连接 + 异步锁延迟创建

RabbitMQ 的 Connection 是重型对象,它底层维护一个 TCP 连接。一个服务实例里通常只需要一个连接,所有 Channel 共享这条连接。但创建连接是异步操作(CreateConnectionAsync()),并发启动时会产生竞态:

private IConnection _connection;          // 单例连接
private readonly SemaphoreSlim _connectionLock = new(1, 1);  // 容量1的异步锁

private async Task<IConnection> GetConnectionAsync()
{
    if (_connection is { IsOpen: true }) return _connection;   // 快速路径

    await _connectionLock.WaitAsync();
    try
    {
        if (_connection is { IsOpen: true }) return _connection;  // 双重检查

        _connectionFactory.HostName = _settings.HostName;
        _connectionFactory.UserName = _settings.UserName;
        _connectionFactory.Password = _settings.Password;
        _connectionFactory.RequestedHeartbeat = TimeSpan.FromSeconds(_settings.RequestedHeartbeat);
        // 随机重连间隔 10~30 秒,防止多实例同时断连时雪崩
        _connectionFactory.NetworkRecoveryInterval = TimeSpan.FromSeconds(new Random().Next(10, 30));

        _connection = await _connectionFactory.CreateConnectionAsync();
        return _connection;
    }
    finally
    {
        _connectionLock.Release();
    }
}

要点: - 双重检查:锁里再判断一次 IsOpen,避免排队的线程重复创建连接 - SemaphoreSlim(1, 1) 而不是 lock:lock 不能 await,在异步代码里必须换成 SemaphoreSlim - 随机重连间隔:NetworkRecoveryInterval 不写死,让多实例同时重连时错开时间

设计二:拓扑缓存,只声明一次

每次发消息都 ExchangeDeclare + QueueDeclare + QueueBind 会让 RabbitMQ 承受不必要的元数据压力。我们用一个 ConcurrentDictionary 做声明幂等:

// 拓扑结构缓存,Key 是 "exchange-queue-routingKey",bool 值本身无意义
private static readonly ConcurrentDictionary<string, bool> _topologyCache
    = new(StringComparer.Ordinal);

private async Task<IChannel> GetChannelAsync(
    string exchange, string queue, string routingKey, string exchangeType, ...)
{
    var con = await GetConnectionAsync();
    var channel = await con.CreateChannelAsync();

    var key = $"{exchange}-{queue}-{routingKey}";
    if (_topologyCache.TryAdd(key, true))   // 只有第一个到达的线程会走进去
    {
        await ExchangeDeclareAsync(channel, exchange, exchangeType, durable, autoDelete, exchangeArgs);
        await QueueDeclareAsync(channel, queue, durable, autoDelete, queueArgs);
        await channel.QueueBindAsync(queue, exchange, routingKey);
    }
    return channel;
}

死信交换机和死信队列也走同样的缓存逻辑,用独立的 key:

if (!string.IsNullOrWhiteSpace(queueInfo.DeadExchangeName))
{
    var key = $"{queueInfo.DeadExchangeName}-{queueInfo.DeadQueueName}-{queueInfo.DeadRoutingKey}";
    if (_topologyCache.TryAdd(key, true))
    {
        await ExchangeDeclareAsync(channel, queueInfo.DeadExchangeName, queueInfo.DeadExchangeTypeName, queueInfo.Durable, false, null);
        await QueueDeclareAsync(channel, queueInfo.DeadQueueName, queueInfo.Durable, false, null);
        await channel.QueueBindAsync(queueInfo.DeadQueueName, queueInfo.DeadExchangeName, queueInfo.DeadRoutingKey);
    }
}

设计三:泛型发送接口 + 反射特性缓存

为什么要用泛型?因为泛型参数 T 直接对应消息模型类,我们可以从 typeof(T) 反射出特性,调用方只需要传消息对象,不需要显式传交换机名或队列名。

private static readonly ConcurrentDictionary<Type, RabbitMqQueueAttribute> _attributeCache = new();

private RabbitMqQueueAttribute GetRabbitMqQueueAttribute<T>()
{
    return _attributeCache.GetOrAdd(typeof(T), t =>
    {
        var attribute = t.GetCustomAttribute<RabbitMqQueueAttribute>();
        if (attribute == null)
            throw new InvalidOperationException($"消息实体 {t.Name} 未配置 RabbitMqQueueAttribute");
        attribute.Validate(t);
        return attribute;
    });
}

GetOrAdd 是线程安全的,配合 ConcurrentDictionary 做到了反射一次、缓存终身——消息模型类上的特性不会被重复读取。

设计四:死信队列参数和延迟消息参数自动组装

RabbitMqClient 把特性配置翻译成 RabbitMQ 原生需要的 Dictionary<string, object> 参数:

// 死信队列参数 → 队列声明时注入
private static IDictionary<string, object> BuildQueueArgs(RabbitMqQueueAttribute queueInfo)
{
    if (string.IsNullOrWhiteSpace(queueInfo.DeadExchangeName)) return null;
    return new Dictionary<string, object>
    {
        {"x-dead-letter-exchange", queueInfo.DeadExchangeName},
        {"x-dead-letter-routing-key", queueInfo.DeadRoutingKey}
    };
}

// 延迟交换机参数 → 交换机声明时注入
private static IDictionary<string, object> BuildExchangeArgs(RabbitMqQueueAttribute queueInfo)
{
    if (queueInfo.ExchangeType != RabbitMqExchangeType.Delayed) return null;
    return new Dictionary<string, object>
    {
        { "x-delayed-type", queueInfo.DelayedExchangeTypeName }
    };
}

// 单条消息的延迟属性 → BasicProperties.Headers 注入
private static BasicProperties BuildProperties(RabbitMqQueueAttribute queueInfo, uint delayMilliseconds)
{
    var properties = new BasicProperties { DeliveryMode = DeliveryModes.Persistent };
    if (queueInfo.ExchangeType == RabbitMqExchangeType.Delayed)
    {
        properties.Headers = new Dictionary<string, object>
        {
            { "x-delay", delayMilliseconds }
        };
    }
    return properties;
}

使用者打特性、客户端自动填参数——这就是声明式配置带来的好处。

设计五:消费者手动 ACK/NACK + 异常兜底

消费者是消息可靠投递的关键。我们强制 autoAck: false,用一个 DoAsync 方法把 ACK/NACK 逻辑收口:

consumer.ReceivedAsync += async (model, ea) =>
{
    await DoAsync(channel, ea, handler, cancellationToken);
};

private async Task DoAsync<T>(IChannel channel, BasicDeliverEventArgs ea,
    Func<T, CancellationToken, Task<bool>> handler, CancellationToken ct)
{
    try
    {
        var msgBody = JsonHelper.Deserialize<T>(Encoding.UTF8.GetString(ea.Body.Span));
        var isSuccess = await handler(msgBody, ct);   // 业务 handler 返回 bool 表达成败

        if (isSuccess)
            await channel.BasicAckAsync(ea.DeliveryTag, multiple: false, cancellationToken: ct);
        else
            await channel.BasicNackAsync(ea.DeliveryTag, multiple: false, requeue: false, cancellationToken: ct);
    }
    catch
    {
        // 业务抛出异常 → 也 NACK 进死信,绝对不能让消息挂在 unacked 状态
        await channel.BasicNackAsync(ea.DeliveryTag, multiple: false, requeue: false, cancellationToken: ct);
    }
}

要点: - handler 返回 Task<bool> 表达成功/失败,比 Action 能携带更多语义 - 异常也走 NACK,任何情况下消息要么 ACK 要么 NACK,绝对不允许 unacked 泄露 - requeue: false:让 NACK 消息进入死信队列,而不是回到原队列重复消费(防止毒消息无限循环)

设计六:QoS prefetch 控制内存

消费者启动时必须设置 BasicQosAsync:

await channel.BasicQosAsync(prefetchSize: 0, prefetchCount: prefetchCount, global: false, cancellationToken: cancellationToken);

prefetchCount 默认 1,意味着未被 ACK 的消息最多 1 条在内存里等待处理。如果 handler 处理慢,RabbitMQ 就不会再推下一条,防止消费者内存溢出。

设计七:IAsyncDisposable + Dispose 桥接

RabbitMQ.Client 6.x 开始全面异步,但有些调用方还是同步代码。我们同时实现 IDisposable 和 IAsyncDisposable,并在 Dispose 里桥接:

public void Dispose()
{
    DisposeAsync().AsTask().GetAwaiter().GetResult();
    GC.SuppressFinalize(this);
}

public async ValueTask DisposeAsync()
{
    // 先取消所有消费者、再关闭通道、最后关闭连接
    foreach (var kv in _consumerChannels)
    {
        try { await kv.Value.BasicCancelAsync(kv.Key); await kv.Value.CloseAsync(); }
        catch { }
    }
    if (_connection != null)
    {
        await _connection.CloseAsync();
        await _connection.DisposeAsync();
        _connection = null;
    }
}

消费者通道没有放在 using 里——因为消费者需要常驻(长连接监听消息),所以我们把它登记到 _consumerChannels 字典,等 Dispose 时统一清理。


四、RabbitMqConfig:自动恢复配置

RabbitMQ 客户端的自动重连是个容易被忽略的点。我们把它做成配置项,默认开启:

public class RabbitMqConfig
{
    public string HostName { get; init; }
    public int Port { get; init; } = 5672;
    public string UserName { get; init; }
    public string Password { get; init; }
    public string VirtualHost { get; init; } = "/";

    /// <summary>连接断开时自动重连,默认开启</summary>
    public bool AutomaticRecoveryEnabled { get; init; } = true;

    /// <summary>重连后自动重建交换机、队列、绑定,默认开启</summary>
    public bool TopologyRecoveryEnabled { get; init; } = true;

    /// <summary>心跳间隔秒数,默认15</summary>
    public int RequestedHeartbeat { get; init; } = 15;
}

在创建连接时赋值到 ConnectionFactory:

_connectionFactory.AutomaticRecoveryEnabled = _settings.AutomaticRecoveryEnabled;
_connectionFactory.TopologyRecoveryEnabled = _settings.TopologyRecoveryEnabled;
_connectionFactory.RequestedHeartbeat = TimeSpan.FromSeconds(_settings.RequestedHeartbeat);

TopologyRecoveryEnabled = true 非常重要——如果没有它,RabbitMQ.Client 在连接恢复后不会重新声明交换机和队列,消息会直接丢失。这也是我们为什么在消费者端也声明一次拓扑的原因(RabbitMQ 官方推荐:生产者和消费者都声明一次,双保险)。


五、使用示例

1. 定义消息模型

[RabbitMqQueue(
    ExchangeName = "switchdata.direct",
    ExchangeType = RabbitMqExchangeType.Direct,
    QueueName = "switchdata.order",
    RoutingKey = "order.created",
    DeadExchangeName = "switchdata.dead.direct",
    DeadQueueName = "switchdata.dead.order",
    DeadRoutingKey = "order.dead")]
public class OrderCreatedEvent
{
    public long OrderId { get; set; }
    public string CustomerId { get; set; }
    public decimal Amount { get; set; }
}

2. 注册客户端到 DI(ASP.NET Core 8)

builder.Services.AddSingleton<RabbitMqClient>(sp =>
{
    var cfg = builder.Configuration.GetSection("RabbitMq").Get<RabbitMqConfig>();
    return new RabbitMqClient(cfg);
});

3. 发送消息(一行代码)

var order = new OrderCreatedEvent
{
    OrderId = 12345, CustomerId = "C001", Amount = 199.00m
};
await rabbitMq.SendMessageAsync(new[] { order });

// 延迟 5 分钟发送
await rabbitMq.SendMessageAsync(new[] { order }, delayMilliseconds: 5 * 60 * 1000);

4. 消费消息

// 可以用 BackgroundService 或 IHostedService 启动消费者
await rabbitMq.ReceiveMessageAsync<OrderCreatedEvent>(async (order, ct) =>
{
    Console.WriteLine($"收到订单 {order.OrderId},金额 {order.Amount}");
    // 返回 true 表示业务成功,消息 ACK 删除
    // 返回 false 或抛异常 → NACK 进死信
    return true;
}, prefetchCount: 1, cancellationToken: stoppingToken);

六、踩过的坑

坑一:TryAdd 的调用时机

ConcurrentDictionary.TryAdd 是原子操作,但它只保证”只有一个线程能成功”,不保证你能拿到锁之后立即执行声明。我们的做法是把 TryAdd 放在锁外、GetChannelAsync 内(连接已经建好之后),这样声明动作一定会在新 Channel 上执行,不会与其他线程冲突。

坑二:消费者通道不能 Dispose

最开始写的时候,发送端和消费端都写成 await using var channel = ...,结果消费者启动几秒后就收不到消息了——因为通道被 Dispose 了。消费者需要长连接常驻,所以我们放弃 using,改用字典登记,在 DisposeAsync 里手动关闭。

坑三:NACK 不能 requeue

RabbitMQ 默认 NACK 的 requeue = true。如果一条消息业务处理永远失败(数据损坏、格式不符),requeue=true 会让它立刻重新入队并再次被消费,形成毒消息死循环。我们统一传 requeue: false,让它进死信队列。

坑四:延迟消息的 ExchangeTypeName

RabbitMQ.Client 的 ExchangeDeclare 方法在指定了 delayed-message-exchange 插件时,类型字符串必须精确匹配插件要求的 x-delayed-message,不能传枚举名 Delayed。我们在特性里做了自动映射,调用方不需要关心。

坑五:BackgroundService 中 Dispose 的顺序

BackgroundService 是 IHostedService,它停止时会传一个 CancellationToken。消费者应该在这个 Token 被取消时停止监听并清理资源,所以我们把 stoppingToken 一路传到底层 channel.BasicConsumeAsync 和 BasicAckAsync,确保优雅下线。


七、总结

这套封装的核心思想可以浓缩成三句话:

  1. 声明式拓扑:消息模型类打一个 [RabbitMqQueue],交换机、队列、路由键、死信、延迟配置全部随模型类走,改拓扑只改这一处
  2. 泛型驱动 + 缓存反射:SendMessageAsync<T> 从 typeof(T) 反射特性,用 ConcurrentDictionary 缓存,做到了反射一次、终身复用
  3. 幂等声明 + 安全消费:拓扑声明靠 ConcurrentDictionary.TryAdd 保证只执行一次,消费端强制手动 ACK/NACK + 异常兜底,绝不允许 unacked 泄露

加上单例连接、自动恢复、异步资源释放这些基础设施,整个封装体虽然不到 400 行,但把 RabbitMQ.Client 的”裸 API”包装成了生产可用的”一键调用”。如果你也在项目里散落地写着 ExchangeDeclare + QueueDeclare 模板代码,不妨试试这种思路——一个特性类 + 一个泛型客户端,就能让消息收发真正做到”业务代码不关心队列细节”。