Files
QYZH.InteractiveMagazine/QYZH.InteractiveMagazine.WorkService/Consumers/RabbitMQHostedService.cs

83 lines
3.0 KiB
C#
Raw Normal View History

using QYZH.InteractiveMagazine.Infrastructure.RabbitMQ;
namespace QYZH.InteractiveMagazine.WorkService.Consumers;
/// <summary>
/// RabbitMQ 消费者后台服务,自动发现并启动所有已注册的 IQueueConsumer
/// </summary>
public class RabbitMQHostedService : BackgroundService
{
private readonly IServiceProvider _serviceProvider;
private readonly ILogger<RabbitMQHostedService> _logger;
public RabbitMQHostedService(IServiceProvider serviceProvider, ILogger<RabbitMQHostedService> logger)
{
_serviceProvider = serviceProvider;
_logger = logger;
}
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
using var scope = _serviceProvider.CreateScope();
var consumers = scope.ServiceProvider.GetServices<IQueueConsumer>().ToList();
if (consumers.Count == 0)
{
_logger.LogWarning("未注册任何队列消费者");
return;
}
_logger.LogInformation("发现 {Count} 个队列消费者,开始启动...", consumers.Count);
var tasks = consumers.Select(c => StartConsumerAsync(c.QueueName, stoppingToken)).ToList();
await Task.WhenAll(tasks);
}
private async Task StartConsumerAsync(string queueName, CancellationToken stoppingToken)
{
_logger.LogInformation("正在启动消费者: {QueueName}", queueName);
try
{
var rabbitMQService = _serviceProvider.GetRequiredService<IRabbitMQService>();
await rabbitMQService.ReceiveAsync(queueName, async (channel, ea) =>
{
// 每条消息创建独立 scope确保消费者内可注入 Scoped 服务
using var messageScope = _serviceProvider.CreateScope();
var consumer = messageScope.ServiceProvider
.GetServices<IQueueConsumer>()
.First(c => c.QueueName == queueName);
try
{
await consumer.HandleAsync(ea.Body.ToArray());
await channel.BasicAckAsync(ea.DeliveryTag, false, stoppingToken);
}
catch (Exception ex)
{
_logger.LogError(ex, "消费者 {QueueName} 处理消息异常", queueName);
await channel.BasicNackAsync(ea.DeliveryTag, false, true, stoppingToken);
try
{
await consumer.OnErrorAsync(ea.Body.ToArray(), ex);
}
catch (Exception errorEx)
{
_logger.LogError(errorEx, "消费者 {QueueName} OnErrorAsync 执行异常", queueName);
}
}
}, stoppingToken);
}
catch (OperationCanceledException)
{
_logger.LogInformation("消费者 {QueueName} 已停止", queueName);
}
catch (Exception ex)
{
_logger.LogError(ex, "消费者 {QueueName} 启动失败", queueName);
}
}
}