消息与后台任务

系统讲解 .NET 后台任务与消息处理体系,覆盖 BackgroundService 与 Channel 队列的应用内解耦、RabbitMQ 与 Azure Service Bus 的集成、消息幂等与重试,以及死信与监控的工程实践。

1. 后台任务与进程内解耦

一句话总结: 后台任务把「请求线程必须等的事」挪到异步执行,BackgroundService 宿主 + 队列缓冲让请求快速返回、任务可靠执行。

Web 请求不该做慢事:发邮件、生成报表、调用第三方、批量写库。把这些放进后台任务,请求立刻返回,任务在后台排队执行。最简单的进程内解耦是「请求 → 内存队列 → BackgroundService 消费」。

public class EmailQueue
{
    private readonly Channel<EmailMessage> _channel =
        Channel.CreateUnbounded<EmailMessage>();

    public void Enqueue(EmailMessage msg) => _channel.Writer.TryWrite(msg);
    public IAsyncEnumerable<EmailMessage> ReadAllAsync(CancellationToken ct) =>
        _channel.Reader.ReadAllAsync(ct);
}

public class EmailSenderWorker(EmailQueue queue, IEmailClient client)
    : BackgroundService
{
    protected override async Task ExecuteAsync(CancellationToken stoppingToken)
    {
        await foreach (var msg in queue.ReadAllAsync(stoppingToken))
        {
            try
            {
                await client.SendAsync(msg);
            }
            catch (Exception ex)
            {
                Console.WriteLine($"发送失败: {msg.To} - {ex.Message}");
            }
        }
    }
}

// 请求侧:投递后立即返回 202
app.MapPost("/contact", (EmailQueue queue, ContactForm form) =>
{
    queue.Enqueue(new EmailMessage(form.Email, "新咨询"));
    return Results.Accepted();
});
队列特点适用
Channel.Unbounded无限缓冲流量波动小
Channel.Bounded有界缓冲 + 背压限流削峰
System.Threading.Channels生产消费模型进程内

避坑: 内存队列的最大局限是进程重启即丢消息——Channel 的消息在内存里,服务重启、崩溃就没了。它只适合「丢了也无所谓」的任务(缓存预热、非关键通知);关键消息必须落到 RabbitMQ / Service Bus 等持久化队列。

2. BackgroundService 与并发消费

一句话总结: 单个 BackgroundService 串行消费是吞吐瓶颈,多实例并行消费 + 有界 Channel 背压,才能让后台任务安全地规模化。

一个 ExecuteAsync 循环一次只处理一个消息,吞吐受限。并行方案有多个 HostedService 实例、单服务内多消费者、或 Parallel.ForEachAsync 控制并发度。关键在并发上限 + 失败隔离——并发不是越多越好。

public class OrderDispatcher : BackgroundService
{
    private readonly Channel<OrderCreatedEvent> _channel;
    private readonly IServiceScopeFactory _scopeFactory;
    private readonly int _concurrency;

    protected override async Task ExecuteAsync(CancellationToken stoppingToken)
    {
        var workers = Enumerable.Range(0, _concurrency)
            .Select(_ => ProcessAsync(stoppingToken));
        await Task.WhenAll(workers);   // 并发消费者
    }

    private async Task ProcessAsync(CancellationToken ct)
    {
        await foreach (var evt in _channel.Reader.ReadAllAsync(ct))
        {
            using var scope = _scopeFactory.CreateScope();
            var handler = scope.ServiceProvider.GetRequiredService<IOrderHandler>();
            try
            {
                await handler.HandleAsync(evt);
            }
            catch (Exception ex)
            {
                Console.WriteLine($"订单处理失败 {evt.OrderId}: {ex.Message}");
            }
        }
    }
}
// 或:Parallel.ForEachAsync 控制并发度
await Parallel.ForEachAsync(
    queue.ReadAllAsync(stoppingToken),
    new ParallelOptions { MaxDegreeOfParallelism = 8, CancellationToken = stoppingToken },
    async (evt, ct) =>
    {
        using var scope = _scopeFactory.CreateScope();
        await scope.ServiceProvider.GetRequiredService<IOrderHandler>()
            .HandleAsync(evt);
    });
并行手段并发度注意
多 HostedService各服务独立需注册多个实例
循环内多消费者固定 N共享 Channel Reader
Parallel.ForEachAsyncMaxDegree 控制异常需自行捕获

避坑: 并行消费把「单点失败」变成「并发失败」——一个消息抛异常可能让 WhenAll 全部退出。每个消息的异常必须就地捕获,或者用「失败入重试队列」而非让消费者整体崩掉。并发数要看下游承受力,别拍脑袋设 100。

3. RabbitMQ 集成

一句话总结: RabbitMQ 用交换机(Exchange)路由到队列(Queue),消费者确认(Ack)与预取(Prefetch)是可靠消费的两个核心旋钮。

RabbitMQ 的模型:生产者发到 Exchange,Exchange 按 routing key 路由到绑定的 Queue,消费者从 Queue 拉取。BasicAck 确认处理完成(否则消息会重新投递),BasicQos(prefetchCount) 控制每个消费者预取数。

// 生产者:发布订单事件
var factory = new ConnectionFactory { HostName = "rabbitmq" };
await using var conn = await factory.CreateConnectionAsync();
await using var channel = await conn.CreateChannelAsync();

await channel.ExchangeDeclareAsync("shop.events",
    type: ExchangeType.Topic, durable: true);
await channel.QueueDeclareAsync("orders.created",
    durable: true, exclusive: false, autoDelete: false);
await channel.QueueBindAsync("orders.created", "shop.events", "order.created");

var body = JsonSerializer.SerializeToUtf8Bytes(new OrderCreatedEvent(orderId));
await channel.BasicPublishAsync(
    exchange: "shop.events",
    routingKey: "order.created",
    body: body,
    mandatory: false);
// 消费者:BackgroundService 内确认
public class OrderCreatedConsumer(IConnectionFactory factory) : BackgroundService
{
    protected override async Task ExecuteAsync(CancellationToken stoppingToken)
    {
        await using var conn = await factory.CreateConnectionAsync();
        await using var channel = await conn.CreateChannelAsync();
        await channel.BasicQosAsync(prefetchCount: 10);   // 预取上限

        var consumer = new AsyncEventingBasicConsumer(channel);
        consumer.ReceivedAsync += async (_, ea) =>
        {
            try
            {
                var evt = JsonSerializer
                    .Deserialize<OrderCreatedEvent>(ea.Body.Span)!;
                await ProcessOrderAsync(evt);
                await channel.BasicAckAsync(ea.DeliveryTag);   // 确认成功
            }
            catch
            {
                await channel.BasicNackAsync(ea.DeliveryTag,
                    requeue: false);   // 进死信,而非无限重投
            }
        };

        await channel.BasicConsumeAsync("orders.created",
            autoAck: false, consumer: consumer);
        await Task.Delay(Timeout.Infinite, stoppingToken);
    }
}
概念作用
Exchange路由中枢(direct/topic/fanout)
Queue消息缓冲与持久化
Ack/Nack确认/否定确认
Prefetch每消费者预取上限
死信队列失败消息的归宿

避坑: autoAck: false 必须显式确认——忘了 Ack 消息会重复投递。处理失败别无限 requeue(消息立即回队头形成死循环),正确姿势是 Nack + requeue:false + 死信队列,把坏消息隔离出来审计。

4. Azure Service Bus 集成

一句话总结: Service Bus 是云托管队列,Queue/Topic 模型与托管语义(Session、DLQ、延迟)开箱即用,连接字符串与命名空间配置即可接入。

Azure Service Bus 分 Queue(单消费者)与 Topic/Subscription(发布订阅)两级。.NET 的 ServiceBusClient + ServiceBusProcessor 提供接收、确认、自动死信。托管在 Azure,无需自建集群,适合已有 Azure 环境的中大型应用。

// 生产者:Session 保证同订单消息有序
var client = new ServiceBusClient(connectionString);
var sender = client.CreateSender("orders.created");

await sender.SendMessageAsync(new ServiceBusMessage(
    JsonSerializer.Serialize(new OrderCreatedEvent(orderId)))
{
    ContentType = "application/json",
    SessionId = orderId.ToString()
});
// 消费者:BackgroundService
public class OrderProcessor(ServiceBusClient client) : BackgroundService
{
    private ServiceBusProcessor? _processor;

    protected override async Task ExecuteAsync(CancellationToken stoppingToken)
    {
        _processor = client.CreateProcessor("orders.created",
            new ServiceBusProcessorOptions
            {
                MaxConcurrentCalls = 8,
                AutoCompleteMessages = false,   // 手动完成
                MaxAutoLockRenewalDuration = TimeSpan.FromMinutes(5)
            });

        _processor.ProcessMessageAsync += async args =>
        {
            try
            {
                var evt = args.Message.Body.ToObjectFromJson<OrderCreatedEvent>();
                await HandleOrderAsync(evt);
                await args.CompleteMessageAsync(args.Message);   // 完成
            }
            catch
            {
                await args.DeadLetterMessageAsync(args.Message,
                    "processing-failed", "处理异常");
            }
        };

        _processor.ProcessErrorAsync += args =>
        {
            Console.WriteLine($"处理器错误: {args.Exception.Message}");
            return Task.CompletedTask;
        };

        await _processor.StartProcessingAsync(stoppingToken);
        await Task.Delay(Timeout.Infinite, stoppingToken);
    }

    public override async Task StopAsync(CancellationToken ct)
    {
        if (_processor is not null)
            await _processor.StopProcessingAsync();
        await base.StopAsync(ct);
    }
}
Service Bus 特性用途
Queue点对点,单消费者
Topic/Subscription发布订阅,多消费者
Session同 key 消息有序
DLQ死信自动隔离
MaxConcurrentCalls并发消费上限

避坑: Service Bus 的「完成」必须和业务成功绑定——AutoCompleteMessages=false 时确认(Complete)前崩溃消息会重新投递。消费处理最好幂等(见第 6 节),配合 Session 保证同订单顺序,避免「确认了但业务没做」与「做了又重投」两种错。

5. 消息重试与退避策略

一句话总结: 瞬时故障(连接抖动、限流)靠指数退避重试,持久故障(数据错误)该进死信;重试要有上限、有间隔、有日志,避免重试风暴。

消息消费失败要区分「可重试」与「不可重试」。可重试的用指数退避(1s → 2s → 4s…)+ 抖动,最多 N 次;不可重试的(反序列化失败、业务规则冲突)直接进死信。盲目无限重试会让故障从「一条消息」蔓延成「重试风暴」。

public class RetryPolicy
{
    private const int MaxAttempts = 5;

    public async Task<bool> TryExecuteAsync(
        Func<Task> action, ILogger log, CancellationToken ct)
    {
        for (int attempt = 1; attempt <= MaxAttempts; attempt++)
        {
            try
            {
                await action();
                return true;
            }
            catch (TransientException ex) when (attempt < MaxAttempts)
            {
                // 指数退避 + 随机抖动
                var delay = TimeSpan.FromMilliseconds(
                    100 * Math.Pow(2, attempt) + Random.Shared.Next(0, 200));
                log.LogWarning("第 {Attempt} 次重试:{Message}", attempt, ex.Message);
                await Task.Delay(delay, ct);
            }
        }
        return false;   // 重试耗尽,交给上层转死信
    }
}
失败类型处理
瞬时(网络、超时、5xx)指数退避重试
业务规则冲突直接失败,不重试
反序列化失败立即进死信
重试耗尽转死信 + 告警

避坑: 两个极端都危险——不重试会让瞬时故障白白丢消息;无限重试会让下游在故障期间雪上加霜。重试要设上限、要退避、要把「重试也算入消息延迟」的意识带进监控;DB 主键冲突、参数非法这类「重试一万次也不会好」的错误,第一时间死信。

6. 消息幂等与去重

一句话总结: 幂等让同一消息处理多次与处理一次效果相同,靠业务唯一键 + 幂等表/状态机实现,是「至少一次投递」语义下的可靠基石。

消息队列的投递语义是「至少一次」——消费者崩溃、确认失败都会导致重投。消费者必须幂等。标准做法:消息带业务唯一键(orderId、eventId),消费时先查幂等表,已处理过直接跳过。

public class OrderHandler(OrderDbContext db, IDistributedCache cache)
{
    public async Task HandleAsync(OrderCreatedEvent evt)
    {
        var dedupKey = $"dedup:order:{evt.OrderId}";

        // 分布式锁 + 已处理标记,双重去重
        await using var lockHandle = await _lock.AcquireAsync(dedupKey, ...);
        if (lockHandle is null) return;               // 他人正在处理
        if (await cache.GetAsync(dedupKey) is not null) return;   // 已处理过
        if (await db.Orders.FindAsync(evt.OrderId) is not null) return;   // 已就位

        db.Orders.Add(new Order { Id = evt.OrderId, ... });
        await db.SaveChangesAsync();

        await cache.SetAsync(dedupKey, [1],
            new DistributedCacheEntryOptions
            { AbsoluteExpirationRelativeToNow = TimeSpan.FromDays(1) });
    }
}
幂等手段原理
业务唯一键 + 查询已存在即跳过
幂等表记录已处理的 messageId
状态机状态已到达目标即跳过
数据库唯一索引重复插入被约束拦截

避坑: 幂等判断和业务写入必须在同一个原子边界里——「先查再写」两个步骤之间有窗口,并发重投会钻空子。数据库唯一索引是最硬的兜底;缓存幂等标记只适合「容忍短暂重复」的场景,且标记要设过期,别让幂等表无限膨胀。

7. 死信、监控与背压

一句话总结: 死信队列承接处理失败的消息,消息监控(积压、延迟、消费速率)与背压机制共同保证消息系统在故障下依旧可控。

消息系统要回答三个问题:卡住的(积压)、失败的(死信)、太快的(背压)。死信队列隔离问题消息;监控指标(队列深度、消费延迟、重试次数)让问题可视化;有界队列 + 限流让生产速率可调节。

// 监控:积压与处理耗时指标
public class QueueMonitor(MessageMetrics metrics) : BackgroundService
{
    protected override async Task ExecuteAsync(CancellationToken stoppingToken)
    {
        using var timer = new PeriodicTimer(TimeSpan.FromSeconds(10));
        while (await timer.WaitForNextTickAsync(stoppingToken))
        {
            var depth = await QueryQueueDepthAsync();   // RabbitMQ 管理 API
            metrics.ObserveQueueDepth("orders.created", depth);
        }
    }
}
监控维度指标告警阈值
积压队列深度持续 > 阈值
延迟消息年龄超出 SLA
消费速率msg/s低于生产速率
失败死信增长率突增即告警
重试重试次数分布接近上限

避坑: 积压不一定是消费慢,也可能是生产风暴或下游变慢——告警要带上下文。背压要双向设计:生产侧有界队列拒绝过快(Channel.CreateBounded 满则等),消费侧 prefetch/concurrency 控制拉取速率。监控的目的是发现「哪个环节是瓶颈」,而不是只看「队列有多深」。

8. 总结

环节要点
后台任务BackgroundService 宿主 + Channel 解耦
并行消费多消费者 + 并发上限 + 失败隔离
RabbitMQExchange 路由 + Ack/Prefetch 可靠消费
Service Bus云托管,Session 保序 + DLQ
重试指数退避 + 上限 + 抖动,持久失败进死信
幂等业务唯一键 + 幂等表 + 唯一索引兜底
监控与背压积压/延迟/死信指标 + 有界队列

消息与后台任务是 .NET 服务从「同步单体」走向「异步分布式」的枢纽:请求只负责接受意图,后台负责执行承诺。这套体系的可靠不在于任何单一组件,而在于投递语义、重试策略、幂等设计与监控告警的组合拳。把「至少一次投递」当作默认现实来设计消费者,把「失败」当作一等公民来设计队列拓扑——消息系统才能从「组件」变成「可靠的地基」。

延伸阅读

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「csharp」更多文章

  1. 缓存与并发控制
  2. 测试体系:xUnit 与 Moq
  3. 托管生命周期与部署