上一篇已经把发布端最小链路搭起来:订单服务通过 orders.events 发布 order.created,库存服务则从 inventory.order-events 接收它。本文把消费端真正落进 ASP.NET Core,重点解决五件事:订阅队列、手动确认、限制未确认消息、区分 prefetch 与并发度,以及在停机时正确排空在途消息。

关于 readyunackedackprefetch 的协议语义,第四篇《RabbitMQ 消费控制:ack、nack、prefetch 与消息确认》已经详细解释过。这里不重复推导,直接关注 RabbitMQ.Client 的代码结构和生命周期。

系列导航:

统一示例:库存服务消费 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);
}

BackgroundServiceAddHostedService 注册为 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。处理顺序是:

  1. 在回调内反序列化消息。
  2. 创建消息级 DI scope,调用库存服务并等待业务完成。
  3. 成功后发送 BasicAckAsync
  4. 只有业务失败才根据错误类型决定是否 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 应结合单条处理耗时、消息体大小、实例数和允许的故障重投规模调整。值太大时,一个慢实例会预占很多消息;值太小时,消费者可能在处理间隙等待下一次投递。81632 只能作为压测起点,不是通用默认值。

需要并发消费时,再显式打开

如果业务可以并行处理,可以把回调并发度调大:

var factory = new ConnectionFactory
{
    // 其他连接参数省略
    ConsumerDispatchConcurrency = 4
};

此时最多有四个消费回调同时执行,而 prefetchCount: 16 仍然只是在途窗口。并发回调可能以不同顺序完成,因此业务不能再依赖“收到顺序等于完成顺序”。

多个回调还会共享同一个消费 Channel。RabbitMQ 官方建议,多线程消费者应对共享 Channel 上的确认操作做互斥。可以把主示例中的 AckAsyncNackAsync 改成下面这种显式传入 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

完整顺序是:

  1. base.StopAsync 取消 stoppingToken,消费会话停止等待新事件。
  2. BasicCancelAsync 通知 Broker 停止后续投递。
  3. 等待 UnregisteredAsync,以 consumer unregistered 作为串行 dispatcher barrier。
  4. 再等待已经进入回调的 InFlight 归零。
  5. 正常关闭 Channel 和 Connection。

只等待 InFlight.IsEmpty 不够,因为客户端内部可能已有 delivery 等待进入回调,计数会在短暂窗口内错误地显示为零。注销事件排在这些 delivery 之后,才构成“不会再有遗漏回调开始”的屏障。这个保证只适用于主示例的单 dispatcher;并发版本必须重新设计任务登记。

如果 _shutdownDeadline 或内部 20 秒期限先到期,代码会释放 Channel/Connection。注意,不是 deadline 本身让消息回队,而是随后发生的 Channel/连接关闭使未确认 delivery 重新入队。宿主的 HostOptions.ShutdownTimeout 还应大于内部排空期限,否则宿主可能先停止等待。

优雅关闭、排空与重投递消费者收到停止信号后取消订阅,等待注销屏障,再把 in-flight 从三条处理到零;若 deadline 到期,则关闭传输并让未确认消息回队后 redelivery。SHUTDOWN PATH · CANCEL FIRST, THEN DRAIN收到停止信号host stoppingBasicCancel停止新投递consumer unregistereddispatcher barrier → 等待已有工作排空IN-FLIGHT WINDOWin-flight3刚取消订阅in-flight2持续完成与 ackin-flight0全部确认MUTUALLY EXCLUSIVE OUTCOMESDRAINED BEFORE DEADLINE处理完成并 ackClose ChannelClose Connection关闭 deadline 到期close transport → requeueQueue · redelivery优雅关闭、排空与重投递移动布局中依次展示停止新投递、consumer unregistered 注销屏障、in-flight 排空,以及成功关闭或传输关闭后回队重投递。GRACEFUL SHUTDOWN收到停止信号BasicCancel · 停止新投递consumer unregistered · 注销屏障in-flight3in-flight2活跃工作仍可走两条互斥结果OUTCOME A · DRAINEDin-flight0处理完成并 ackClose Channel → Close Connection有限时间内顺利收尾OUTCOME B · DEADLINE / REQUEUE关闭 deadline 到期close transport → requeueredelivery之后由 Queue 再次投递Queueredelivery回队路径避开结果标题
优雅关闭、排空与重投递注销屏障防止漏掉客户端内已派发的 delivery;有限时间内未确认的消息要等 Channel 或连接关闭后才会由 Broker 重新入队。
消费者先取消订阅,等待 consumer unregistered 形成 dispatcher barrier,再排空 in-flight;超时路径是 deadline → close transport → requeue,而不是 deadline 直接让消息回队。

优雅关闭的目标不是承诺“当前消息一定处理完”,而是在明确的时间预算内尽量排空;超过预算后,关闭传输并依靠 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 放进同一个 catchack 失败被误判成业务失败,可能错误 nack分开处理业务失败与确认结果未知
所有异常都 requeue: true永久错误形成高速重复投递使用有限重试并把不可恢复消息送入 DLX
只处理连接后的自动恢复初次连接失败或 Channel 级异常后消费者消失选定一种恢复所有权,并覆盖会话创建与终止
停机收尾继续使用已取消的 stoppingToken取消订阅和排空立即失败用独立的关闭期限令牌执行收尾
认为 redelivery 只会发生在业务异常时ack 前断线造成的重复无法解释把未确认重投递视为正常故障语义

本篇完成了消费端的基础实现:消息级 DI scope 隔离业务依赖,manual ack 定义业务完成边界,prefetch 控制在途窗口,应用循环统一重建失败会话,StopAsync 则触发标准停止,再由注销屏障和关闭期限完成有限排空。

下一篇:RabbitMQ 发布可靠性:mandatory、publisher confirms 与不可路由消息(九)