Files
glz cbdee5068a feat: 新增消息Outbox机制、雪花ID配置优化及多项功能完善
1.  新增数据库唯一约束和Message_Outbox表脚本
2.  新增雪花ID、Hangfire存储、MQ重试等配置实体
3.  重构各项目雪花ID生成逻辑,改为从配置读取WorkerId
4.  优化积分服务分页查询、用户背包更新逻辑
5.  新增JWT令牌Redis过期刷新逻辑
6.  完善RabbitMQ死信队列消息头信息
7.  新增可靠MQ消息发布服务和Outbox派发后台服务
8.  替换原有RabbitMQ直接发送为Outbox可靠发布
9.  优化签到服务逻辑,新增重复签到校验和补签卡扣减逻辑
10. 修复自动铺码消费逻辑,新增点阵页预占和释放机制
2026-07-10 10:44:00 +08:00

135 lines
5.7 KiB
C#
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

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, stoppingToken)).ToList();
await Task.WhenAll(tasks);
}
private async Task StartConsumerAsync(IQueueConsumer queueConsumer, CancellationToken stoppingToken)
{
var queueName = queueConsumer.QueueName;
var exchange = queueConsumer.Exchange;
var routingKey = queueConsumer.RoutingKey;
_logger.LogInformation("正在启动消费者: {QueueName}, Exchange: {Exchange}, RoutingKey: {RoutingKey}", queueName, exchange, routingKey);
try
{
var rabbitMQService = _serviceProvider.GetRequiredService<IRabbitMQService>();
await rabbitMQService.ReceiveAsync(exchange, queueName, routingKey, async (channel, ea) =>
{
// 声明死信队列及绑定幂等操作确保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);
// 每条消息创建独立 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(), stoppingToken);
await channel.BasicAckAsync(ea.DeliveryTag, false, stoppingToken);
}
catch (Exception ex)
{
_logger.LogError(ex, "消费者 {QueueName} 处理消息异常", queueName);
// 发送到死信队列,成功则从主队列移除,失败则重新入队
var dlqSent = await SendToDeadLetterQueueAsync(exchange, dlqName, dlqRoutingKey, queueName, ea.Body.ToArray(), ex, stoppingToken);
await channel.BasicNackAsync(ea.DeliveryTag, false, !dlqSent, stoppingToken);
if (dlqSent)
{
_logger.LogInformation("消息已转入死信队列 {DlqName}", dlqName);
}
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);
}
}
/// <summary>
/// 将失败消息发送到死信队列
/// </summary>
private async Task<bool> SendToDeadLetterQueueAsync(string exchange, string dlqName, string dlqRoutingKey, string originalQueueName, byte[] body, Exception exception, 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,
Headers = new Dictionary<string, object?>
{
["x-original-queue"] = originalQueueName,
["x-error-type"] = exception.GetType().FullName,
["x-error-message"] = exception.Message,
["x-failed-at"] = DateTimeOffset.UtcNow.ToString("O")
}
};
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;
}
}
}