在日常 .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,确保优雅下线。
七、总结
这套封装的核心思想可以浓缩成三句话:
- 声明式拓扑:消息模型类打一个
[RabbitMqQueue],交换机、队列、路由键、死信、延迟配置全部随模型类走,改拓扑只改这一处 - 泛型驱动 +
缓存反射:
SendMessageAsync<T>从typeof(T)反射特性,用ConcurrentDictionary缓存,做到了反射一次、终身复用 - 幂等声明 + 安全消费:拓扑声明靠
ConcurrentDictionary.TryAdd保证只执行一次,消费端强制手动 ACK/NACK + 异常兜底,绝不允许 unacked 泄露
加上单例连接、自动恢复、异步资源释放这些基础设施,整个封装体虽然不到
400 行,但把 RabbitMQ.Client 的”裸
API”包装成了生产可用的”一键调用”。如果你也在项目里散落地写着
ExchangeDeclare + QueueDeclare
模板代码,不妨试试这种思路——一个特性类 +
一个泛型客户端,就能让消息收发真正做到”业务代码不关心队列细节”。