RabbitMQ + .NET 实战:消费消息、手动确认与优雅关闭(八)
上一篇已经把发布端最小链路搭起来:订单服务通过 orders.events 发布 order.created,库存服务则从 inventory.order-events 接收它。本文把消费端真正落进 ASP.NET Core,重点解决五件事:订阅队列、手动确认、限制未确认消息、区分 prefetch 与并发度,以及在停机时正确排空在途消息。
关于 ready、unacked、ack 和 prefetch 的协议语义,第四篇《RabbitMQ 消费控制:ack、nack、prefetch 与消息确认》已经详细解释过。这里不重复推导,直接关注 RabbitMQ.Client 的代码结构和生命周期。
系列导航:
- 上一篇:RabbitMQ + .NET 实战:建立连接、声明拓扑与发布消息(七)
- 当前篇:RabbitMQ + .NET 实战:消费消息、手动确认与优雅关闭(八)(本文)
- 下一篇:RabbitMQ 发布可靠性:mandatory、publisher confirms 与不可路由消息(九)
统一示例:库存服务消费 order.created
沿用上一篇的消息:
{
"eventId": "01J7X8A6YQ4Q9M6Q4QG7H2J8F1",
"eventType": "order.created",
"occurredAt": "2026-07-15T10:00:00Z",
"orderId": "ORD-20260715-001",
"customerId": "C10086",
"amount": 199.00,
"currency": "CNY"
}
C# 属性保持 PascalCase,通过 JsonPropertyName 显式对应消息中的 camelCase 字段:
using System.Text.Json.Serialization;
public sealed record OrderCreated(
[property: JsonPropertyName("eventId")] string EventId,
[property: JsonPropertyName("eventType")] string EventType,
[property: JsonPropertyName("occurredAt")] DateTimeOffset OccurredAt,
[property: JsonPropertyName("orderId")] string OrderId,
[property: JsonPropertyName("customerId")] string CustomerId,
[property: JsonPropertyName("amount")] decimal Amount,
[property: JsonPropertyName("currency")] string Currency);
库存预留属于业务代码,可以抽象成 scoped 依赖。接口从一开始就接收 eventId,让后续幂等实现不必再改变消费边界:
public interface IInventoryReservationService
{
Task ReserveAsync(
string eventId,
string orderId,
CancellationToken cancellationToken);
}
BackgroundService 由 AddHostedService 注册为 singleton,不能直接捕获通常依赖 scoped DbContext 的业务服务。因此消费者只注入 IServiceScopeFactory,每条投递单独创建作用域。Microsoft 的后台服务作用域指南也明确指出,Hosted Service 不会自动获得请求作用域。
先看完整的 BackgroundService
下面是一套连贯的基础实现。它默认使用单回调串行消费,并明确采用“应用循环重建整个消费会话”的恢复策略:初次连接、Connection 或 Channel 结束后都回到同一个带退避的重试入口,不再同时叠加客户端自动恢复。
using System.Collections.Concurrent;
using System.Text.Json;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging;
using RabbitMQ.Client;
using RabbitMQ.Client.Events;
public sealed class InventoryConsumerService(
IServiceScopeFactory scopeFactory,
ILogger<InventoryConsumerService> logger) : BackgroundService
{
private static readonly TimeSpan SessionDrainTimeout =
TimeSpan.FromSeconds(20);
// StopAsync 把宿主的“停止优雅关闭”期限转发给会话清理逻辑。
private readonly CancellationTokenSource _shutdownDeadline = new();
private sealed class ConsumerSessionState
{
public ConcurrentDictionary<Guid, byte> InFlight { get; } = new();
public TaskCompletionSource Unregistered { get; } = new(
TaskCreationOptions.RunContinuationsAsynchronously);
public TaskCompletionSource<Exception> Faulted { get; } = new(
TaskCreationOptions.RunContinuationsAsynchronously);
}
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
var factory = new ConnectionFactory
{
HostName = "localhost",
Port = 5672,
UserName = "blog_demo",
Password = "local-demo-change-me",
VirtualHost = "/",
ClientProvidedName = "inventory-worker",
// RabbitMQ.Client 默认就是 1,这里显式写出消费模型。
ConsumerDispatchConcurrency = 1,
// 本文选择由外层循环统一重建 Connection、Channel 和 Consumer。
AutomaticRecoveryEnabled = false
};
var retryDelay = TimeSpan.FromSeconds(1);
while (!stoppingToken.IsCancellationRequested)
{
IConnection? connection = null;
IChannel? channel = null;
string? consumerTag = null;
var state = new ConsumerSessionState();
try
{
connection = await factory.CreateConnectionAsync(stoppingToken);
channel = await connection.CreateChannelAsync(
cancellationToken: stoppingToken);
await channel.BasicQosAsync(
prefetchSize: 0,
prefetchCount: 16,
global: false,
cancellationToken: stoppingToken);
var channelStopped = new TaskCompletionSource<ShutdownEventArgs>(
TaskCreationOptions.RunContinuationsAsynchronously);
channel.ChannelShutdownAsync += (_, reason) =>
{
channelStopped.TrySetResult(reason);
return Task.CompletedTask;
};
var consumer = new AsyncEventingBasicConsumer(channel);
consumer.UnregisteredAsync += (_, _) =>
{
state.Unregistered.TrySetResult();
return Task.CompletedTask;
};
consumer.ReceivedAsync += (_, args) =>
HandleMessageAsync(channel, state, args);
consumerTag = await channel.BasicConsumeAsync(
queue: "inventory.order-events",
autoAck: false,
consumer: consumer,
cancellationToken: stoppingToken);
// 一旦成功订阅,下一次故障从最短退避重新开始。
retryDelay = TimeSpan.FromSeconds(1);
Task completed = await Task.WhenAny(
channelStopped.Task,
state.Unregistered.Task,
state.Faulted.Task)
.WaitAsync(stoppingToken);
if (completed == state.Faulted.Task)
{
Exception callbackException = await state.Faulted.Task;
throw new InvalidOperationException(
"Consumer callback lost its safe disposition boundary.",
callbackException);
}
if (completed == state.Unregistered.Task)
{
throw new InvalidOperationException(
"Consumer was cancelled unexpectedly.");
}
ShutdownEventArgs reason = await channelStopped.Task;
throw new InvalidOperationException(
$"Consumer channel closed: {reason.ReplyCode} {reason.ReplyText}");
}
catch (OperationCanceledException)
when (stoppingToken.IsCancellationRequested)
{
// 进入 finally,用独立的关闭期限取消订阅并排空。
}
catch (Exception exception)
{
logger.LogWarning(
exception,
"RabbitMQ consumer session ended; reconnecting in {Delay}",
retryDelay);
}
finally
{
await DrainAndCloseSessionAsync(
connection,
channel,
consumerTag,
state);
}
if (stoppingToken.IsCancellationRequested)
{
break;
}
var jitter = TimeSpan.FromMilliseconds(Random.Shared.Next(0, 500));
try
{
await Task.Delay(retryDelay + jitter, stoppingToken);
}
catch (OperationCanceledException)
when (stoppingToken.IsCancellationRequested)
{
break;
}
retryDelay = TimeSpan.FromSeconds(
Math.Min(retryDelay.TotalSeconds * 2, 30));
}
}
private async Task HandleMessageAsync(
IChannel channel,
ConsumerSessionState state,
BasicDeliverEventArgs args)
{
var workId = Guid.NewGuid();
state.InFlight.TryAdd(workId, 0);
try
{
OrderCreated message;
try
{
// Body 只能在回调返回前读取、反序列化或复制。
message = JsonSerializer.Deserialize<OrderCreated>(args.Body.Span)
?? throw new JsonException("Message body is null.");
}
catch (JsonException exception)
{
logger.LogError(exception, "Invalid order.created payload");
await NackAsync(
channel,
args.DeliveryTag,
requeue: false);
return;
}
try
{
await using var scope = scopeFactory.CreateAsyncScope();
var inventory = scope.ServiceProvider
.GetRequiredService<IInventoryReservationService>();
// 业务处理使用自己的超时,不复用宿主的停止令牌。
using var processingCts = new CancellationTokenSource(
TimeSpan.FromSeconds(30));
await inventory.ReserveAsync(
message.EventId,
message.OrderId,
processingCts.Token);
}
catch (Exception businessException)
{
logger.LogError(
businessException,
"Failed to reserve inventory; requeue = {Requeue}",
!args.Redelivered);
// 仅用于演示:首次失败立即重回队列,redelivery 再失败则拒绝。
await NackAsync(
channel,
args.DeliveryTag,
requeue: !args.Redelivered);
return;
}
try
{
await AckAsync(channel, args.DeliveryTag);
}
catch (Exception ackException)
{
// 业务已完成,不能再把确认异常伪装成业务失败并发送 nack。
logger.LogError(
ackException,
"Inventory reserved but ack outcome is unknown; " +
"EventId = {EventId}",
message.EventId);
// 结束当前会话并关闭 Channel,让未确认消息按协议重新投递。
state.Faulted.TrySetResult(ackException);
}
}
catch (Exception dispositionException)
{
logger.LogError(
dispositionException,
"Consumer could not reach a safe ack/nack disposition");
state.Faulted.TrySetResult(dispositionException);
}
finally
{
state.InFlight.TryRemove(workId, out _);
}
}
private static async Task AckAsync(
IChannel channel,
ulong deliveryTag)
{
await channel.BasicAckAsync(
deliveryTag: deliveryTag,
multiple: false);
}
private static async Task NackAsync(
IChannel channel,
ulong deliveryTag,
bool requeue)
{
await channel.BasicNackAsync(
deliveryTag: deliveryTag,
multiple: false,
requeue: requeue);
}
private async Task DrainAndCloseSessionAsync(
IConnection? connection,
IChannel? channel,
string? consumerTag,
ConsumerSessionState state)
{
using var cleanupCts = CancellationTokenSource.CreateLinkedTokenSource(
_shutdownDeadline.Token);
cleanupCts.CancelAfter(SessionDrainTimeout);
CancellationToken cleanupToken = cleanupCts.Token;
try
{
if (channel is { IsOpen: true } && consumerTag is not null)
{
if (!state.Unregistered.Task.IsCompleted)
{
await channel.BasicCancelAsync(
consumerTag: consumerTag,
cancellationToken: cleanupToken);
}
// cancel-ok 只停止后续投递;等待 Unregistered 形成 dispatcher barrier。
await state.Unregistered.Task.WaitAsync(cleanupToken);
}
while (!state.InFlight.IsEmpty)
{
await Task.Delay(50, cleanupToken);
}
if (channel is { IsOpen: true })
{
await channel.CloseAsync(cancellationToken: cleanupToken);
}
if (connection is { IsOpen: true })
{
await connection.CloseAsync(cancellationToken: cleanupToken);
}
}
catch (OperationCanceledException)
when (cleanupToken.IsCancellationRequested)
{
logger.LogWarning(
"RabbitMQ consumer shutdown deadline elapsed; " +
"closing transport so unacknowledged deliveries can be redelivered");
}
catch (Exception exception)
{
logger.LogWarning(exception, "RabbitMQ consumer session close failed");
}
finally
{
await DisposeSessionAsync(channel, connection);
}
}
private async Task DisposeSessionAsync(
IChannel? channel,
IConnection? connection)
{
try
{
if (channel is not null)
{
await channel.DisposeAsync();
}
}
catch (Exception disposeException)
{
logger.LogDebug(
disposeException,
"RabbitMQ consumer session dispose failed for Channel");
}
try
{
if (connection is not null)
{
await connection.DisposeAsync();
}
}
catch (Exception disposeException)
{
logger.LogDebug(
disposeException,
"RabbitMQ consumer session dispose failed for Connection");
}
}
public override async Task StopAsync(CancellationToken shutdownToken)
{
using var registration = shutdownToken.Register(
static state => ((CancellationTokenSource)state!).Cancel(),
_shutdownDeadline);
// base.StopAsync 立即取消 stoppingToken,ExecuteAsync 再进入会话排空。
await base.StopAsync(shutdownToken);
}
}
最后把消费者注册进依赖注入:
builder.Services.AddScoped<IInventoryReservationService, InventoryReservationService>();
builder.Services.AddHostedService<InventoryConsumerService>();
builder.Services.Configure<HostOptions>(options =>
{
// 必须大于示例内部 20 秒的排空期限,给宿主留出退出余量。
options.ShutdownTimeout = TimeSpan.FromSeconds(30);
});
真实项目中,连接参数、排空期限和重试上限都应来自配置,凭据则来自密钥管理。这里保留固定值,只为了让消费流程集中在一个示例中。
这段代码显式设置 AutomaticRecoveryEnabled = false,不是说自动恢复不可用,而是避免“客户端自动恢复”和“应用重建会话”同时争夺生命周期。当前策略覆盖两类自动恢复本身不会替你兜底的情况:初次连接失败不会触发自动恢复,Channel 级异常也不会自动恢复关闭的 Channel。若项目选择 RabbitMQ.Client 自动恢复,就应删除外层重建策略,并另外处理这两个边界;不要把两套恢复机制混在一起。具体限制见官方 .NET 客户端恢复说明。
消息为什么必须在业务完成后确认
示例固定使用 autoAck: false。处理顺序是:
- 在回调内反序列化消息。
- 创建消息级 DI scope,调用库存服务并等待业务完成。
- 成功后发送
BasicAckAsync。 - 只有业务失败才根据错误类型决定是否 nack;ack 异常属于“业务已完成、确认结果未知”。
如果使用 autoAck: true,RabbitMQ 在投递消息后就可以认为这次消费已经完成。此时进程即使在库存事务提交前崩溃,Broker 也无法再提供这次处理机会。对允许丢失的低价值通知,自动确认可能是合理取舍;对订单、库存和支付等关键业务,通常应使用手动确认。
这里特意把业务 try/catch 和确认 try/catch 分开。库存事务已经提交后,如果 BasicAckAsync 因连接中断而失败,消费者不能再记录“库存预留失败”,更不能根据 Redelivered 发送 requeue: false。示例会结束当前会话并关闭传输,让未确认消息重新投递;下一篇再用 EventId 幂等收敛这段结果未知窗口。
args.Body 是 RabbitMQ.Client 管理的 ReadOnlyMemory<byte>。回调返回后不能继续持有它;需要交给回调外任务时,必须先调用 ToArray() 复制。示例在回调内立即反序列化,因此不需要额外复制。这个约束可以在 RabbitMQ .NET 客户端指南 中核对。
nack 不是完整的重试系统
完整示例采用了一条刻意简化的失败规则:
- JSON 无法解析:
requeue: false,因为同一份错误数据不会靠立即重试自愈。 - 业务首次失败:
requeue: true,允许一次立即重新投递。 args.Redelivered已经为true:再次失败后使用requeue: false。
这只能避免最直接的无限热循环,不能代替生产级有限重试。Redelivered 只是“这次投递以前是否投递过”的标记,不是可靠的重试计数器。
还要特别注意:使用 requeue: false 时,只有队列配置了 Dead Letter Exchange,消息才会进入死信链路;没有配置 DLX,消息会被直接丢弃。如何区分瞬时失败、永久失败和毒消息,应继续阅读《RabbitMQ 消费可靠性:幂等、重试与死信队列》;不要把本文的简化分支直接当成最终生产策略。
prefetch 控制在途窗口,不控制回调并发度
示例中的两个数字属于不同层次:
ConsumerDispatchConcurrency = 1 → 本进程同时执行几个消费回调
prefetchCount = 16 → 这个消费者最多持有多少条未确认消息
因此,主示例虽然设置了 prefetchCount: 16,但回调仍然默认串行执行。多出来的消息只是处于 unacked 并等待前面的回调结束,不会自动变成 16 个并行业务任务。
prefetch 应结合单条处理耗时、消息体大小、实例数和允许的故障重投规模调整。值太大时,一个慢实例会预占很多消息;值太小时,消费者可能在处理间隙等待下一次投递。8、16 或 32 只能作为压测起点,不是通用默认值。
需要并发消费时,再显式打开
如果业务可以并行处理,可以把回调并发度调大:
var factory = new ConnectionFactory
{
// 其他连接参数省略
ConsumerDispatchConcurrency = 4
};
此时最多有四个消费回调同时执行,而 prefetchCount: 16 仍然只是在途窗口。并发回调可能以不同顺序完成,因此业务不能再依赖“收到顺序等于完成顺序”。
多个回调还会共享同一个消费 Channel。RabbitMQ 官方建议,多线程消费者应对共享 Channel 上的确认操作做互斥。可以把主示例中的 AckAsync 和 NackAsync 改成下面这种显式传入 Channel 的形式:
private readonly SemaphoreSlim _channelGate = new(1, 1);
private async Task AckAsync(IChannel channel, ulong deliveryTag)
{
await _channelGate.WaitAsync();
try
{
await channel.BasicAckAsync(
deliveryTag: deliveryTag,
multiple: false);
}
finally
{
_channelGate.Release();
}
}
private async Task NackAsync(
IChannel channel,
ulong deliveryTag,
bool requeue)
{
await _channelGate.WaitAsync();
try
{
await channel.BasicNackAsync(
deliveryTag: deliveryTag,
multiple: false,
requeue: requeue);
}
finally
{
_channelGate.Release();
}
}
并发确认仍然保持 multiple: false,避免不同完成顺序下重复确认一段 delivery tag。但只加锁还不够:主示例的“先等 UnregisteredAsync,再检查活动回调”依赖 ConsumerDispatchConcurrency = 1 的串行 dispatcher。若提升并发度,应使用应用自有的有界工作队列或任务登记屏障,保证 cancel-ok 之前已经派发的每个回调都被纳入排空范围,不能继续用某一瞬间的空计数判断完成。关于默认顺序分发、ConsumerDispatchConcurrency 和共享 Channel 确认约束,可以参考 RabbitMQ .NET 客户端并发说明。
优雅关闭:停止信号与关闭期限各司其职
ASP.NET Core 的两个令牌表达不同含义:
| 令牌 | 表达的含义 | 在本文中的用途 |
|---|---|---|
ExecuteAsync(stoppingToken) | 后台任务应该开始结束 | 终止消费会话等待并进入 finally |
StopAsync(shutdownToken) | 宿主还允许优雅关闭多久 | 限制取消订阅、排空与关闭传输的总时间 |
BackgroundService.StopAsync 的标准行为是先取消 stoppingToken,再等待 ExecuteAsync 退出。因此示例先注册宿主的 shutdownToken,随后立即调用 base.StopAsync(shutdownToken):消费循环收到停止信号后进入 finally,清理逻辑则使用尚未取消的 _shutdownDeadline。
完整顺序是:
base.StopAsync取消stoppingToken,消费会话停止等待新事件。BasicCancelAsync通知 Broker 停止后续投递。- 等待
UnregisteredAsync,以consumer unregistered作为串行 dispatcher barrier。 - 再等待已经进入回调的
InFlight归零。 - 正常关闭 Channel 和 Connection。
只等待 InFlight.IsEmpty 不够,因为客户端内部可能已有 delivery 等待进入回调,计数会在短暂窗口内错误地显示为零。注销事件排在这些 delivery 之后,才构成“不会再有遗漏回调开始”的屏障。这个保证只适用于主示例的单 dispatcher;并发版本必须重新设计任务登记。
如果 _shutdownDeadline 或内部 20 秒期限先到期,代码会释放 Channel/Connection。注意,不是 deadline 本身让消息回队,而是随后发生的 Channel/连接关闭使未确认 delivery 重新入队。宿主的 HostOptions.ShutdownTimeout 还应大于内部排空期限,否则宿主可能先停止等待。
优雅关闭的目标不是承诺“当前消息一定处理完”,而是在明确的时间预算内尽量排空;超过预算后,关闭传输并依靠 manual ack 的重投递语义恢复。ASP.NET Core 对停止令牌和关闭期限的定义可以参考 BackgroundService 文档。
连接丢失后,应该预期 redelivery
只要使用 manual ack,就必须接受一个现实:连接丢失、Channel 关闭、进程崩溃或宿主重启之后,尚未确认的投递可能再次出现。
如果 ReserveAsync() 已经提交数据库,但连接恰好在 BasicAckAsync() 到达 Broker 前断开,那么 RabbitMQ 无法知道业务是否完成,只能重新投递。这不是 Broker 出错,而是 at-least-once 模型在故障窗口里的正常结果。
因此,消费端必须接受两个事实:
- “业务代码执行过”不等于 Broker 已经收到确认。
- redelivery 可能发生在任意一次正常业务处理之后,而不只是代码抛出异常时。
- ack 失败不能反推出业务失败,也不能据此安全地改发 nack。
本文只建立这个预期。下一阶段要用稳定的 EventId、业务事务和幂等记录把重复结果收敛掉,而不能试图用 prefetch、并发度或优雅关闭消除所有重复。
消费端常见误区
| 误区 | 直接后果 | 更稳妥的做法 |
|---|---|---|
把 autoAck: true 理解成“业务成功后自动确认” | 业务执行前消息就可能失去重投机会 | 关键业务通常先处理,再手动 ack |
| 在 singleton Hosted Service 中直接注入 scoped 业务服务 | DbContext 生命周期错误,或被提升为跨消息共享实例 | 每条消息创建 DI scope |
| 把 prefetch 当成并发度 | 调大后只增加 unacked,业务仍可能串行 | 分开配置在途窗口与回调并发度 |
| 并发回调直接共享确认 Channel | 确认操作可能并发交错 | 对 ack/nack 做互斥并保持 multiple: false |
| 把业务处理和 ack 放进同一个 catch | ack 失败被误判成业务失败,可能错误 nack | 分开处理业务失败与确认结果未知 |
所有异常都 requeue: true | 永久错误形成高速重复投递 | 使用有限重试并把不可恢复消息送入 DLX |
| 只处理连接后的自动恢复 | 初次连接失败或 Channel 级异常后消费者消失 | 选定一种恢复所有权,并覆盖会话创建与终止 |
停机收尾继续使用已取消的 stoppingToken | 取消订阅和排空立即失败 | 用独立的关闭期限令牌执行收尾 |
| 认为 redelivery 只会发生在业务异常时 | ack 前断线造成的重复无法解释 | 把未确认重投递视为正常故障语义 |
本篇完成了消费端的基础实现:消息级 DI scope 隔离业务依赖,manual ack 定义业务完成边界,prefetch 控制在途窗口,应用循环统一重建失败会话,StopAsync 则触发标准停止,再由注销屏障和关闭期限完成有限排空。