2026-06-11 18:01:07 +08:00
|
|
|
|
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);
|
|
|
|
|
|
|
2026-06-22 16:56:05 +08:00
|
|
|
|
var tasks = consumers.Select(c => StartConsumerAsync(c, stoppingToken)).ToList();
|
2026-06-11 18:01:07 +08:00
|
|
|
|
await Task.WhenAll(tasks);
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-06-22 16:56:05 +08:00
|
|
|
|
private async Task StartConsumerAsync(IQueueConsumer queueConsumer, CancellationToken stoppingToken)
|
2026-06-11 18:01:07 +08:00
|
|
|
|
{
|
2026-06-22 16:56:05 +08:00
|
|
|
|
var queueName = queueConsumer.QueueName;
|
|
|
|
|
|
var exchange = queueConsumer.Exchange;
|
|
|
|
|
|
var routingKey = queueConsumer.RoutingKey;
|
|
|
|
|
|
_logger.LogInformation("正在启动消费者: {QueueName}, Exchange: {Exchange}, RoutingKey: {RoutingKey}", queueName, exchange, routingKey);
|
2026-06-11 18:01:07 +08:00
|
|
|
|
|
|
|
|
|
|
try
|
|
|
|
|
|
{
|
|
|
|
|
|
var rabbitMQService = _serviceProvider.GetRequiredService<IRabbitMQService>();
|
|
|
|
|
|
|
2026-06-22 16:56:05 +08:00
|
|
|
|
await rabbitMQService.ReceiveAsync(exchange, queueName, routingKey, async (channel, ea) =>
|
2026-06-11 18:01:07 +08:00
|
|
|
|
{
|
2026-06-22 16:56:05 +08:00
|
|
|
|
// 声明死信队列及绑定(幂等操作,确保DLQ存在)
|
|
|
|
|
|
var dlqName = $"{queueName}.dlq";
|
|
|
|
|
|
var dlqRoutingKey = $"{routingKey}.dlq";
|
|
|
|
|
|
await channel.QueueDeclareAsync(queue: dlqName, durable: true, exclusive: false, autoDelete: false, arguments: null);
|
|
|
|
|
|
await channel.QueueBindAsync(queue: dlqName, exchange: exchange, routingKey: dlqRoutingKey, arguments: null);
|
|
|
|
|
|
|
2026-06-11 18:01:07 +08:00
|
|
|
|
// 每条消息创建独立 scope,确保消费者内可注入 Scoped 服务
|
|
|
|
|
|
using var messageScope = _serviceProvider.CreateScope();
|
|
|
|
|
|
var consumer = messageScope.ServiceProvider
|
|
|
|
|
|
.GetServices<IQueueConsumer>()
|
|
|
|
|
|
.First(c => c.QueueName == queueName);
|
|
|
|
|
|
|
|
|
|
|
|
try
|
|
|
|
|
|
{
|
2026-06-23 16:53:18 +08:00
|
|
|
|
await consumer.HandleAsync(ea.Body.ToArray(), stoppingToken);
|
2026-06-11 18:01:07 +08:00
|
|
|
|
await channel.BasicAckAsync(ea.DeliveryTag, false, stoppingToken);
|
|
|
|
|
|
}
|
|
|
|
|
|
catch (Exception ex)
|
|
|
|
|
|
{
|
|
|
|
|
|
_logger.LogError(ex, "消费者 {QueueName} 处理消息异常", queueName);
|
2026-06-22 16:56:05 +08:00
|
|
|
|
|
|
|
|
|
|
// 发送到死信队列,成功则从主队列移除,失败则重新入队
|
|
|
|
|
|
var dlqSent = await SendToDeadLetterQueueAsync(exchange, dlqName, dlqRoutingKey, ea.Body.ToArray(), stoppingToken);
|
|
|
|
|
|
await channel.BasicNackAsync(ea.DeliveryTag, false, !dlqSent, stoppingToken);
|
|
|
|
|
|
|
|
|
|
|
|
if (dlqSent)
|
|
|
|
|
|
{
|
|
|
|
|
|
_logger.LogInformation("消息已转入死信队列 {DlqName}", dlqName);
|
|
|
|
|
|
}
|
2026-06-11 18:01:07 +08:00
|
|
|
|
|
|
|
|
|
|
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);
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
2026-06-22 16:56:05 +08:00
|
|
|
|
|
|
|
|
|
|
/// <summary>
|
|
|
|
|
|
/// 将失败消息发送到死信队列
|
|
|
|
|
|
/// </summary>
|
|
|
|
|
|
private async Task<bool> SendToDeadLetterQueueAsync(string exchange, string dlqName, string dlqRoutingKey, byte[] body, CancellationToken cancellationToken)
|
|
|
|
|
|
{
|
|
|
|
|
|
try
|
|
|
|
|
|
{
|
|
|
|
|
|
var rabbitMQConnection = _serviceProvider.GetRequiredService<IRabbitMQConnection>();
|
|
|
|
|
|
using var channel = await rabbitMQConnection.CreateChannel();
|
|
|
|
|
|
await channel.QueueDeclareAsync(queue: dlqName, durable: true, exclusive: false, autoDelete: false, arguments: null);
|
|
|
|
|
|
await channel.QueueBindAsync(queue: dlqName, exchange: exchange, routingKey: dlqRoutingKey, arguments: null);
|
|
|
|
|
|
|
|
|
|
|
|
var properties = new RabbitMQ.Client.BasicProperties
|
|
|
|
|
|
{
|
|
|
|
|
|
Persistent = true
|
|
|
|
|
|
};
|
|
|
|
|
|
await channel.BasicPublishAsync(exchange, dlqRoutingKey, false, properties, body, cancellationToken);
|
|
|
|
|
|
|
|
|
|
|
|
_logger.LogInformation("消息已发送到死信队列: {DlqName}", dlqName);
|
|
|
|
|
|
return true;
|
|
|
|
|
|
}
|
|
|
|
|
|
catch (Exception ex)
|
|
|
|
|
|
{
|
|
|
|
|
|
_logger.LogError(ex, "发送消息到死信队列 {DlqName} 失败", dlqName);
|
|
|
|
|
|
return false;
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
2026-06-11 18:01:07 +08:00
|
|
|
|
}
|