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.ForEachAsync | MaxDegree 控制 | 异常需自行捕获 |
避坑: 并行消费把「单点失败」变成「并发失败」——一个消息抛异常可能让
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 解耦 |
| 并行消费 | 多消费者 + 并发上限 + 失败隔离 |
| RabbitMQ | Exchange 路由 + Ack/Prefetch 可靠消费 |
| Service Bus | 云托管,Session 保序 + DLQ |
| 重试 | 指数退避 + 上限 + 抖动,持久失败进死信 |
| 幂等 | 业务唯一键 + 幂等表 + 唯一索引兜底 |
| 监控与背压 | 积压/延迟/死信指标 + 有界队列 |
消息与后台任务是 .NET 服务从「同步单体」走向「异步分布式」的枢纽:请求只负责接受意图,后台负责执行承诺。这套体系的可靠不在于任何单一组件,而在于投递语义、重试策略、幂等设计与监控告警的组合拳。把「至少一次投递」当作默认现实来设计消费者,把「失败」当作一等公民来设计队列拓扑——消息系统才能从「组件」变成「可靠的地基」。
延伸阅读
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。