RabbitMQ + .NET 10 实战:可靠地发布与消费消息
上一篇建立了 Producer、Exchange、Binding、Queue、Consumer、ack 和 prefetch 的心智模型。本篇把这条链路落实为一组 .NET 10 代码:启动本地 RabbitMQ,使用 RabbitMQ.Client 7.2.1 声明拓扑,可靠地发布 order.created,再以手动确认方式消费它。
本文只解决客户端与 Broker 之间的基础可靠性。数据库与 RabbitMQ 双写、消费者幂等、有限重试、Quorum Queue 和集群不在本篇范围内;这些问题需要端到端设计,不能靠一个客户端参数顺带解决。
系列导航:
- RabbitMQ 核心原理:一条消息如何完成路由与消费
- RabbitMQ + .NET 10 实战:可靠地发布与消费消息(本文)
- RabbitMQ 生产实践:可靠性、幂等、重试与高可用
🐇 启动最小 RabbitMQ 环境
先用 Docker Compose 启动一个带管理界面的单节点 RabbitMQ。下面的账号和口令只供本地演示:
services:
rabbitmq:
image: rabbitmq:4.3.2-management
container_name: rabbitmq-blog-demo
environment:
RABBITMQ_DEFAULT_USER: blog_demo
RABBITMQ_DEFAULT_PASS: local-demo-change-me
ports:
- "5672:5672"
- "15672:15672"
volumes:
- rabbitmq-data:/var/lib/rabbitmq
volumes:
rabbitmq-data:
在 Compose 文件所在目录依次执行:
docker compose up -d
docker compose ps
docker compose down
启动后访问 http://localhost:15672,可以检查节点状态、Exchange、Queue 和消息速率;应用连接使用的 AMQP 地址是 localhost:5672。演示账号、口令以及直接暴露的端口都不能原样用于生产环境,生产部署还需要最小权限、密钥管理、TLS 和网络隔离。
示例目标框架为 .NET 10,直接安装官方客户端,不引入 MassTransit、NServiceBus 或 EasyNetQ:
dotnet add package RabbitMQ.Client --version 7.2.1
🔌 复用 Connection,按职责管理 Channel
为突出生命周期,下面省略配置绑定、日志与密钥管理:
using RabbitMQ.Client;
var factory = new ConnectionFactory
{
HostName = "localhost",
Port = 5672,
UserName = "blog_demo",
Password = "local-demo-change-me",
ClientProvidedName = "orders-api-publisher",
AutomaticRecoveryEnabled = true
};
await using IConnection connection =
await factory.CreateConnectionAsync(cancellationToken);
var channelOptions = new CreateChannelOptions(
publisherConfirmationsEnabled: true,
publisherConfirmationTrackingEnabled: true);
await using IChannel channel =
await connection.CreateChannelAsync(channelOptions, cancellationToken);
Connection 对应一条经过 TCP 建连、协议握手和认证的昂贵连接,应在应用生命周期内长时间复用,绝不能为每条消息新建。Channel 是复用在 Connection 上的轻量协议通道,可以按发布、消费或其他职责分配,生命周期通常短于 Connection。
复用不等于任意并发共享。同一个 IChannel 不应被多个发布线程无保护地同时使用,否则可能造成协议帧交错以及发布确认归属风险。并发发布时,应按职责或并发单元分配 Channel;确实共享时,必须在应用层串行化访问。后面的消费片段沿用变量名以便阅读,真实服务通常会为发布和消费建立各自的 Channel。
AutomaticRecoveryEnabled 能在网络故障后恢复连接、Channel 和已知拓扑,却不会替应用缓存断线期间尝试发布的消息,也不会消除发布结果不确定性。可靠发布仍然需要 publisher confirms 和应用自己的失败处理。
🧱 声明可重复执行的拓扑
订单事件发送到 Topic Exchange orders.events,库存服务使用 Queue inventory.order-events 接收 order.*:
const string exchange = "orders.events";
const string queue = "inventory.order-events";
await channel.ExchangeDeclareAsync(
exchange: exchange,
type: ExchangeType.Topic,
durable: true,
autoDelete: false,
arguments: null,
cancellationToken: cancellationToken);
await channel.QueueDeclareAsync(
queue: queue,
durable: true,
exclusive: false,
autoDelete: false,
arguments: null,
cancellationToken: cancellationToken);
await channel.QueueBindAsync(
queue: queue,
exchange: exchange,
routingKey: "order.*",
arguments: null,
cancellationToken: cancellationToken);
这些声明可以在应用启动时重复执行:对象不存在就创建,存在且属性一致就保持不变。但“同名”并不足够,同名 Exchange 或 Queue 的类型、durable、exclusive、auto-delete 等属性必须与已有对象一致。不一致会触发通道级协议异常并关闭当前 Channel,应用不能继续在这个已关闭的 Channel 上发布或消费。
📤 发布带属性的订单事件
消息除了 JSON 正文,还应携带可供诊断与消费者识别的标准属性。发布前注册 BasicReturnAsync,让 mandatory: true 的不可路由结果不再静默丢失:
using System.Text.Json;
using RabbitMQ.Client.Events;
var orderCreated = new
{
eventId = "01J...",
eventType = "order.created",
occurredAt = DateTimeOffset.Parse("2026-07-15T10:00:00Z"),
orderId = "ORD-20260715-001",
customerId = "C10086"
};
channel.BasicReturnAsync += (_, args) =>
{
Console.Error.WriteLine(
$"Unroutable: {args.ReplyCode} {args.ReplyText}, routingKey={args.RoutingKey}");
return Task.CompletedTask;
};
ReadOnlyMemory<byte> body = JsonSerializer.SerializeToUtf8Bytes(orderCreated);
var properties = new BasicProperties
{
ContentType = "application/json",
DeliveryMode = DeliveryModes.Persistent,
MessageId = orderCreated.eventId,
Type = orderCreated.eventType,
Timestamp = new AmqpTimestamp(orderCreated.occurredAt.ToUnixTimeSeconds())
};
await channel.BasicPublishAsync(
exchange: "orders.events",
routingKey: "order.created",
mandatory: true,
basicProperties: properties,
body: body,
cancellationToken: cancellationToken);
这里有三层不同的保证,不能把它们混成一句“消息发送成功”:
DeliveryMode = DeliveryModes.Persistent表示消息需要持久化。它只有与 durable Queue 等拓扑配置配合,才构成 Broker 重启后的磁盘恢复链条;单独把消息设为 persistent 并不完整。mandatory: true要求消息至少路由到一个 Queue。无法匹配 Binding 时,Broker 会发送basic.return;真实应用应把BasicReturnAsync接入指标、告警和发布失败策略,而不是只写控制台。- 创建 Channel 时同时启用了 publisher confirmations 和客户端确认跟踪。RabbitMQ.Client 7.2.1 会为发布建立关联,等待
BasicPublishAsync就是在等待 Broker 的 ack;Broker nack 或 mandatory return 会使该调用以发布异常失败。因此调用方仍须捕获并记录发布异常,决定重试、持久化待发记录或中止业务。
Publisher confirm 只表示 Broker 已接管这次发布,不表示消费者已经收到消息,更不表示库存业务或数据库事务已经完成。网络也可能在 Broker 接收消息后、确认抵达客户端前中断,留下“Broker 已接收,但客户端没有收到确认”的不确定状态。此时重发可能产生重复,后续必须依靠消息 ID 和消费者幂等来收敛。
| 阶段 | 机制 | 能证明什么 | 不能证明什么 |
|---|---|---|---|
| Publisher → Broker | Publisher confirms | Broker 已接管发布 | 消息一定路由、业务一定成功 |
| Exchange → Queue | mandatory + return | 不可路由时可被发布者观察 | 消费者已经处理 |
| Broker → Consumer | 手动 ack | Consumer 声明本次处理完成 | 数据库与 ack 之间不存在故障窗口 |
📥 异步消费、有限 prefetch 与手动确认
下面假设事件契约定义为:
public sealed record OrderCreated(
string eventId,
string eventType,
DateTimeOffset occurredAt,
string orderId,
string customerId);
ProcessOrderCreatedAsync、异常类型、受控重试基础设施、inFlightMessages 和 consumerChannelShutdown 都是示意接口,本文不提供仓储或监督器实现。代码先把每个消费者允许的未确认消息数限制为 16,再注册手动确认消费者:
using System.Text.Json;
using RabbitMQ.Client.Events;
await channel.BasicQosAsync(
prefetchSize: 0,
prefetchCount: 16,
global: false,
cancellationToken: cancellationToken);
using var processingCts = new CancellationTokenSource();
CancellationToken processingToken = processingCts.Token; // 应用拥有,不是 stoppingToken
var consumer = new AsyncEventingBasicConsumer(channel);
consumer.ReceivedAsync += async (_, args) =>
{
var inFlight = inFlightMessages.Register();
try
{
// RabbitMQ.Client 复用消息体内存;必须在回调返回前复制或完成反序列化。
OrderCreated message;
try
{
message = JsonSerializer.Deserialize<OrderCreated>(args.Body.Span)
?? throw new JsonException("Message body is null.");
}
catch (JsonException)
{
await channel.BasicNackAsync(
deliveryTag: args.DeliveryTag,
multiple: false,
requeue: false,
cancellationToken: processingToken);
return;
}
try
{
await ProcessOrderCreatedAsync(message, processingToken);
await channel.BasicAckAsync(
deliveryTag: args.DeliveryTag,
multiple: false,
cancellationToken: processingToken);
}
catch (TransientDependencyException)
{
await SendToControlledRetryAsync(message, processingToken);
await channel.BasicAckAsync(
deliveryTag: args.DeliveryTag,
multiple: false,
cancellationToken: processingToken);
}
}
catch (Exception unexpected)
{
Console.Error.WriteLine($"Unexpected consumer failure: {unexpected}");
// 监督器必须幂等、立即返回,并在回调外 abort 当前消费 Channel。
consumerChannelShutdown.RequestAbort(channel, unexpected);
}
finally
{
inFlight.Dispose();
}
};
string consumerTag = await channel.BasicConsumeAsync(
queue: "inventory.order-events",
autoAck: false,
consumer: consumer,
cancellationToken: cancellationToken);
RabbitMQ.Client 7.x 以库拥有的 ReadOnlyMemory<byte> 传递消息体;ReceivedAsync 返回后,这段内存可能被回收或复用。因此必须像示例这样在回调内完成反序列化。如果要把原始字节交给回调外的异步任务,先调用 args.Body.ToArray() 复制,再保存副本。
解码与业务处理使用两个独立的 try。只有 JSON 解码或反序列化失败会进入 requeue: false 分支;ProcessOrderCreatedAsync 内部即使抛出 JsonException,也不会被误判成坏消息。业务层应使用明确的领域异常参与重试或停车分类。未分类异常不会 ack 或 nack 当前消息,而是先记录错误,再请求应用监督器在回调外 abort 这个消费专用 Channel;Channel 结束后,其中尚未确认的消息会重新入队。监督器必须让请求幂等并负责按服务策略重建 Channel 和 Consumer,不能把未知异常改成无界的 requeue: true。
RabbitMQ.Client 7.2.1 的异步消费者调度器会捕获逃出回调的异常,并通过 Channel 的 CallbackExceptionAsync 上报。这个事件适合接入日志和告警,但它只是观测机制:不会主动关闭 Channel 或 Connection,也不会让未确认消息自动重投递。若应用既不 ack/nack,也不执行上述受控关闭,消息会一直停留在 unacknowledged 状态,并可能逐步占满 prefetch。示例让监督器在回调外执行 abort,是因为 CloseAsync 与 AbortAsync 都会等待 Channel shutdown 完成;不要在当前消费回调内直接等待自己的 Channel 关闭。
临时错误分支表达的是“先可靠转入受控重试,再 ack 原消息”。第三篇才会设计重试 Exchange、延迟和次数上限。如果 SendToControlledRetryAsync 失败,后面的 ack 不会执行,原消息仍保持未确认;不能捕获该失败后继续 ack,否则会同时失去原消息和重试副本。
永久反序列化错误使用 requeue: false,避免同一毒消息原地循环。消息会被丢弃还是进入死信 Exchange(DLX),取决于 inventory.order-events 的队列策略,本篇没有暗示 DLX 已经配置。
prefetch 不是越大越好
prefetchCount: 16 只是便于观察的演示起点,不是推荐生产值。实际调优至少要同时考虑:
- 单条消息处理耗时;
- 消费者实例数;
- 回调并发度;
- 数据库连接池;
- 第三方 API 限额;
- 允许的未确认消息内存。
提高 prefetch 可能减少消费者等待、提升吞吐,也会增加 unacknowledged 消息数量、故障后的重投递批量以及单个消费者对消息的占用。尤其当回调并发度大于 1 时,prefetch 只是待处理上限之一,不能替代应用自己的并发限制和下游容量保护。
🛑 在 BackgroundService 中优雅关闭
ASP.NET Core BackgroundService 收到停止信号后,应按固定顺序收尾:先停止接收新消息,再等待已经进入处理流程的任务,然后关闭 Channel,最后关闭 Connection。消费回调不能直接捕获 ExecuteAsync 的 stoppingToken:它在停机开始时立即取消,会让已经开始的业务处理、重试转存或 ack 提前中断。上面的 processingToken 由应用拥有,只在 drain deadline 到期后才取消。
using var drainDeadline = new CancellationTokenSource(drainTimeout);
CancellationToken drainDeadlineToken = drainDeadline.Token;
bool drainCompleted = false;
try
{
await channel.BasicCancelAsync(
consumerTag,
cancellationToken: drainDeadlineToken);
await inFlightMessages.WaitForAllAsync(drainDeadlineToken);
drainCompleted = true;
}
finally
{
if (!drainCompleted)
{
// 取消消费或等待在途失败/超时:停止处理,让未 ack 消息随连接关闭重新入队。
processingCts.Cancel();
}
// 清理绝不能复用可能已经取消的 drainDeadlineToken。
using var channelCleanupDeadline =
new CancellationTokenSource(cleanupTimeout);
try
{
await channel.CloseAsync(
cancellationToken: channelCleanupDeadline.Token);
}
catch (Exception closeError)
{
Console.Error.WriteLine($"Channel close failed: {closeError.Message}");
await channel.AbortAsync(CancellationToken.None);
}
finally
{
// Connection 使用全新的清理预算,不能沿用已取消的 Channel 清理令牌。
using var connectionCleanupDeadline =
new CancellationTokenSource(cleanupTimeout);
try
{
await connection.CloseAsync(
cancellationToken: connectionCleanupDeadline.Token);
}
catch (Exception closeError)
{
Console.Error.WriteLine($"Connection close failed: {closeError.Message}");
await connection.AbortAsync(CancellationToken.None);
}
}
}
drainTimeout 是应用配置的在途处理预算,drainDeadlineToken 只控制取消订阅和等待在途,不会因为 BackgroundService 刚收到 stopping signal 就立即取消。只要 BasicCancelAsync、等待在途或 drain deadline 任一环节没有正常完成,finally 就会取消 processingCts,再关闭连接,使未 ack 消息尽快重新入队。
资源清理不能复用可能已经取消的 drainDeadlineToken。示例用 cleanupTimeout 为 Channel 和 Connection 分别创建全新的清理期限;一个清理令牌超时不会阻止下一个资源的关闭尝试。正常 CloseAsync 失败时,再用未取消的 CancellationToken.None 调用 AbortAsync 兜底。真实服务应同时记录这些关闭错误,交给宿主的总停机期限处理。
inFlightMessages 是应用层维护的接口,不是 RabbitMQ.Client 自动提供的能力。它需要像消费示例那样在回调入口登记,并在回调退出时统一注销;此时 ack、nack、受控重试转存后的 ack,以及请求异常关闭 Channel 等终态都已经完成或被明确触发。停止流程应能等待当前快照;BasicCancelAsync 只是取消订阅、阻止新的投递,不代表业务任务已经完成。
consumerChannelShutdown 同样属于应用层监督器。未知异常路径请求的是异常关闭当前消费 Channel,正常停机仍按上面的顺序先 BasicCancelAsync、等待 inFlightMessages,再关闭 Channel 和 Connection;两条路径都应幂等处理“Channel 已在关闭”的竞态。
Channel 和 Connection 同时实现了异步释放;示例外围的 await using 会负责释放资源,显式 CloseAsync 则让正常关闭顺序和超时边界更清晰。若停机超时后进程被终止,尚未 ack 的消息会在连接关闭后重新入队,并可能投递给其他消费者;这也是消费者必须支持重复投递的原因。
⚠️ 常见错误对照
| 错误做法 | 直接后果 | 正确方向 |
|---|---|---|
| 每条消息创建 Connection | 握手、认证与 TCP 成本放大 | 长时间复用 Connection |
autoAck: true 后执行业务 | 业务失败时消息已被删除 | 成功后手动 ack |
所有异常都 requeue: true | 毒消息原地无限循环 | 分类错误、有限重试、最终停车 |
| 只把消息设为 persistent | 队列或发布阶段仍有缺口 | durable 拓扑 + persistent + confirms |
| 多线程无保护共享发布 Channel | 帧交错和确认归属风险 | 按职责分配或串行化 Channel |
下一篇:把客户端保证连成端到端可靠性
到这里,代码已经覆盖长生命周期连接、幂等拓扑声明、publisher confirms、mandatory return、手动 ack、有限 prefetch 和优雅关闭。但数据库提交与消息发布之间仍有双写窗口;业务成功与 ack 之间仍会产生重复;临时失败也还没有次数上限和最终停车位置。
下一篇将围绕同一个 order.created 事件,继续解决 Outbox、消费者幂等、有限重试、Quorum Queue 和故障演练,把本篇的客户端机制组织成可验证的生产方案。
系列导航:
- RabbitMQ 核心原理:一条消息如何完成路由与消费
- RabbitMQ + .NET 10 实战:可靠地发布与消费消息(本文)
- RabbitMQ 生产实践:可靠性、幂等、重试与高可用