RabbitMQ + .NET 实战:建立连接、声明拓扑与发布消息(七)
从这篇开始,系列进入 dotnet-practice 组。目标不是一次讲完所有 RabbitMQ 能力,而是围绕同一个 order.created 事件,把生产者和消费者分别落成能运行、能解释、能扩展的 .NET 10 代码。本文只讲发布端最先要搭好的六件事:ConnectionFactory、连接与 Channel 生命周期、并发发布时的 Channel 所有权、Exchange/Queue/Binding 声明、首次发布,以及发布时值得设置的消息属性。
本文刻意不展开 publisher confirms、mandatory、returned message、消费者手动确认、prefetch 和优雅关闭;这些内容分别留给后两篇,避免把“先发出第一条可靠的消息”和“完整可靠性设计”混成一篇。
系列导航:
- 上一篇:RabbitMQ 常见模式:Work Queue、Pub/Sub、Routing、RPC(六)
- 当前篇:RabbitMQ + .NET 实战:建立连接、声明拓扑与发布消息(七)(本文)
- 下一篇:RabbitMQ + .NET 实战:消费消息、手动确认与优雅关闭(八)
统一示例:order.created
三篇都使用同一个订单事件,避免示例在发布端、消费端和可靠性讨论之间反复切换:
{
"eventId": "01J7X8A6YQ4Q9M6Q4QG7H2J8F1",
"eventType": "order.created",
"occurredAt": "2026-07-15T10:00:00Z",
"orderId": "ORD-20260715-001",
"customerId": "C10086",
"amount": 199.00,
"currency": "CNY"
}
在本组文章里,订单服务把它发布到 Topic Exchange orders.events;库存服务消费 inventory.order-events;通知服务消费 notification.order-events。本文先只把发布端链路搭起来。
ConnectionFactory 先决定连接边界
.NET 客户端的起点几乎总是 ConnectionFactory。它不是一次性的配置桶,而是“以后创建出来的连接应当具备哪些网络和协议属性”的模板。最小示例如下:
using RabbitMQ.Client;
var factory = new ConnectionFactory
{
HostName = "localhost",
Port = 5672,
UserName = "blog_demo",
Password = "local-demo-change-me",
VirtualHost = "/",
ClientProvidedName = "orders-api-publisher"
};
这几个属性里,值得先建立明确心智的是:
HostName、Port、UserName、Password、VirtualHost定义“连到哪个 Broker 上的哪个逻辑空间”。ClientProvidedName用于在管理界面、日志和监控中识别连接来源,线上非常值得设置。ConnectionFactory可以长期复用;它本身不代表已经连上 RabbitMQ,真正昂贵的是后面创建出来的Connection。
如果应用已经有统一配置系统,应该把这些参数从 appsettings.json 或密钥系统读入,而不是硬编码在代码里。本文直接写死,只是为了聚焦生命周期。
连接长寿命,Channel 按职责分配
RabbitMQ 中最容易写错的第一步,是把连接和 Channel 都当成“随手 new 一个、发完就丢”的轻量对象。正确的理解刚好相反:
- Connection 是客户端到 Broker 的长生命周期 TCP 连接,包含握手、认证和协议协商成本,应该在应用生命周期内尽量复用。
- Channel 是复用在 Connection 之上的轻量 AMQP 通道,声明拓扑、发布消息都在这里发生,通常比 Connection 更适合按职责划分。
最小发布端代码如下:
using RabbitMQ.Client;
var factory = new ConnectionFactory
{
HostName = "localhost",
Port = 5672,
UserName = "blog_demo",
Password = "local-demo-change-me",
VirtualHost = "/",
ClientProvidedName = "orders-api-publisher"
};
await using IConnection connection =
await factory.CreateConnectionAsync(cancellationToken);
await using IChannel channel =
await connection.CreateChannelAsync(cancellationToken: cancellationToken);
这里先记住两个生命周期原则:
- 不要为每条消息创建一个新连接。那会把 TCP 建连、认证和资源占用直接放大成吞吐瓶颈。
- 不要把一个
IChannel无保护地交给多个并发发布线程共享。Channel 是轻量对象,但不是“天然线程安全的并发发布总线”。
更实际的做法通常是:
- 每个进程维护少量长寿命连接。
- 按职责分配 Channel,例如一个发布 Channel、一个消费 Channel。
- 确实存在高并发发布时,再按并发单元分配多个 Channel,而不是让所有任务争用同一个 Channel。
本文先建立“一个连接 + 一个发布 Channel”的最小路径;一旦发布动作会被多个请求同时触发,就必须继续明确这个 Channel 的所有权。
并发发布时,先给 IChannel 明确所有权
“多个请求可以同时发布消息”不等于“多个请求可以同时操作同一个 IChannel”。一条 AMQP 消息会经过 method、header、body 等协议帧;如果两个执行流无保护地往同一个 Channel 写入,帧可能交错,Broker 看到的就不再是两条完整消息。RabbitMQ 的 .NET 客户端把这视为发布端的硬性约束:并发发布任务不能共享同一个 IChannel;需要共享时,应用必须自己做互斥。官方并发说明
最简单的做法是让一个 SemaphoreSlim 保护共享 Channel:
var publishGate = new SemaphoreSlim(1, 1);
async Task PublishAsync(
IChannel channel,
ReadOnlyMemory<byte> body,
CancellationToken cancellationToken)
{
await publishGate.WaitAsync(cancellationToken);
try
{
await channel.BasicPublishAsync(
exchange: "orders.events",
routingKey: "order.created",
body: body,
cancellationToken: cancellationToken);
}
finally
{
publishGate.Release();
}
}
这保证了共享 Channel 的正确性,但代价也很直接:所有 Publish 都会串行。真正需要并行时,不是移除锁,而是让每个并发单元在使用期间独占不同的 IChannel。这些 Channel 可以复用同一个长寿命 Connection;Channel Pool 也应当是“借出一个独占 Channel”,而不是把池中的某个 Channel 再次共享给多个发布线程。
| 模型 | Publish 并行度 | 适合什么场景 | 它不解决什么 |
|---|---|---|---|
共享 IChannel + SemaphoreSlim | 无 | 低频发布、实现优先简单 | 请求与发布耗时耦合、进程崩溃后的消息恢复 |
多个独占 IChannel / 有界 Pool | 有限并行 | 发布吞吐确实成为瓶颈 | 数据库与 RabbitMQ 的双写窗口 |
Channel<T> + 单后台 Publisher | 后台串行 | 需要在进程内排队和背压 | 内存消息持久化、多个实例协调 |
第三种模型中的 Channel<T> 是 System.Threading.Channels 的进程内队列,不是 RabbitMQ 的 AMQP Channel。它如何提供背压、又为什么不能替代可靠消息存储,见顺序与背压(五)。如果消息属于订单、支付等不能漏发的业务事件,还需要后续的 Transactional Outbox 把“待发布”写入本地事务。
先把拓扑声明成可重复执行的启动步骤
为了让 order.created 能路由到库存和通知两个业务,我们先声明一个 Topic Exchange,再声明两个 Queue,并分别建立 Binding:
订单服务 Producer
│ order.created
▼
Topic Exchange: orders.events
│
├─ order.* ──> Queue: inventory.order-events
└─ order.# ──> Queue: notification.order-events
对应代码如下:
using RabbitMQ.Client;
const string exchangeName = "orders.events";
const string inventoryQueue = "inventory.order-events";
const string notificationQueue = "notification.order-events";
await channel.ExchangeDeclareAsync(
exchange: exchangeName,
type: ExchangeType.Topic,
durable: true,
autoDelete: false,
arguments: null,
cancellationToken: cancellationToken);
await channel.QueueDeclareAsync(
queue: inventoryQueue,
durable: true,
exclusive: false,
autoDelete: false,
arguments: null,
cancellationToken: cancellationToken);
await channel.QueueDeclareAsync(
queue: notificationQueue,
durable: true,
exclusive: false,
autoDelete: false,
arguments: null,
cancellationToken: cancellationToken);
await channel.QueueBindAsync(
queue: inventoryQueue,
exchange: exchangeName,
routingKey: "order.*",
arguments: null,
cancellationToken: cancellationToken);
await channel.QueueBindAsync(
queue: notificationQueue,
exchange: exchangeName,
routingKey: "order.#",
arguments: null,
cancellationToken: cancellationToken);
这段代码有三个关键点:
- 它适合作为应用启动步骤的一部分重复执行。对象不存在时创建;已经存在且属性一致时保持不变。
- “名字一样”还不够。若同名 Queue 或 Exchange 的类型、
durable、exclusive、autoDelete等属性不一致,Broker 会抛出协议错误并关闭当前 Channel。 - 声明拓扑是为了让发布依赖有显式边界。不要把“RabbitMQ 里总会有人提前建好 Exchange 和 Queue”当成默认前提。
在团队协作里,常见做法是由基础设施代码或应用启动代码负责声明;只要职责明确即可。最怕的是一部分环境靠手工点管理台创建,另一部分环境靠应用自动声明,最后谁也说不清线上真实拓扑是什么。
第一次发布 order.created
拓扑具备后,就可以发布这条消息。先定义一个清晰、稳定的消息载荷:
using System.Text.Json;
var orderCreated = new
{
eventId = "01J7X8A6YQ4Q9M6Q4QG7H2J8F1",
eventType = "order.created",
occurredAt = DateTimeOffset.Parse("2026-07-15T10:00:00Z"),
orderId = "ORD-20260715-001",
customerId = "C10086",
amount = 199.00m,
currency = "CNY"
};
ReadOnlyMemory<byte> body = JsonSerializer.SerializeToUtf8Bytes(orderCreated);
然后调用 BasicPublishAsync:
await channel.BasicPublishAsync(
exchange: "orders.events",
routingKey: "order.created",
body: body,
cancellationToken: cancellationToken);
这就是 RabbitMQ 发布端的最小骨架:
- 用
ConnectionFactory定义连接配置。 - 创建长寿命
Connection。 - 从连接创建发布
Channel。 - 声明 Exchange、Queue 和 Binding。
- 以
routingKey = "order.created"发布 JSON 消息。
做到这里,已经足够让一条消息进入正确拓扑。但“能发出去”和“发布结果可证明”还不是一回事。本文先停在“最小发布路径”;发布确认、不可路由检测与结果不确定窗口,留到第 9 篇单独展开。
发布时值得设置哪些属性
虽然本篇不讲 publisher confirms,但第一次发布时就应该养成设置消息属性的习惯。真正有价值的不是“把所有字段都填满”,而是挑出那些能帮助消费者识别、帮助运维排查、帮助后续幂等设计的属性。
var properties = new BasicProperties
{
ContentType = "application/json",
ContentEncoding = "utf-8",
DeliveryMode = DeliveryModes.Persistent,
MessageId = orderCreated.eventId,
Type = orderCreated.eventType,
Timestamp = new AmqpTimestamp(orderCreated.occurredAt.ToUnixTimeSeconds()),
AppId = "orders-api"
};
await channel.BasicPublishAsync(
exchange: "orders.events",
routingKey: "order.created",
basicProperties: properties,
body: body,
cancellationToken: cancellationToken);
这些属性分别解决不同问题:
ContentType = "application/json":告诉下游按 JSON 解释消息体。ContentEncoding = "utf-8":让编码约定显式化,而不是靠双方猜测。DeliveryMode = Persistent:表达“这是一条需要持久化的消息”;它要和 durable 拓扑一起看,单独设置并不自动代表完整可靠性。MessageId:给这条事件一个稳定 ID,后续做重发、去重和问题排查时都很关键。Type:直接表明业务事件类型,例如order.created。Timestamp:保留事件发生时间,有利于判断延迟和排查积压。AppId:标识消息来源应用,便于多发布者场景下定位来源。
如果事件契约长期演进,还可以再加:
CorrelationId:需要把当前消息和上游请求、流程实例或 saga 关联起来时使用。Headers:放少量稳定的元数据,例如租户、版本或追踪上下文;不要把大块业务正文塞进 headers。
一个实用判断标准是:消费者真正用来执行业务的数据放在消息体里;用于诊断、路由补充或跨系统追踪的稳定元数据,才考虑放到属性或 headers 里。
发布端常见误区
| 误区 | 直接后果 | 更稳妥的做法 |
|---|---|---|
| 每发一条消息就创建一次 Connection | TCP 和认证成本被放大 | 连接长时间复用 |
| 一个 Channel 被多个线程无保护并发使用 | 协议帧交错、连接或客户端异常 | 同一 Channel 串行访问,或让并发单元独占不同 Channel |
| 只发消息,不声明拓扑 | 环境漂移、依赖边界不清 | 启动时显式声明 Exchange/Queue/Binding |
不设置 MessageId | 后续重发、排查、去重都缺少稳定标识 | 事件创建时就生成稳定 ID |
把 DeliveryMode = Persistent 当成“已经绝对可靠” | 高估 Broker 接管边界 | 把它视为可靠性链条中的一环 |
本篇解决了什么:
- 建立了
.NET发布端的最小实现路径:从ConnectionFactory、长寿命连接与Channel分工,到可重复执行的 Exchange、Queue、Binding 声明。 - 明确了并发发布的 Channel 所有权:共享
IChannel必须串行化;需要并行时使用多个独占 Channel,而不是移除并发保护。 - 完成了
order.created的首次发布,并梳理了发布时最值得设置的消息属性,为后续消费端实践和可靠性话题打下基础。