using QYZH.InteractiveMagazine.Infrastructure.RabbitMQ; namespace QYZH.InteractiveMagazine.WorkService.Consumers; /// /// RabbitMQ 消费者后台服务,自动发现并启动所有已注册的 IQueueConsumer /// public class RabbitMQHostedService : BackgroundService { private readonly IServiceProvider _serviceProvider; private readonly ILogger _logger; public RabbitMQHostedService(IServiceProvider serviceProvider, ILogger logger) { _serviceProvider = serviceProvider; _logger = logger; } protected override async Task ExecuteAsync(CancellationToken stoppingToken) { using var scope = _serviceProvider.CreateScope(); var consumers = scope.ServiceProvider.GetServices().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(); await rabbitMQService.ReceiveAsync(queueName, async (channel, ea) => { // 每条消息创建独立 scope,确保消费者内可注入 Scoped 服务 using var messageScope = _serviceProvider.CreateScope(); var consumer = messageScope.ServiceProvider .GetServices() .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); } } }