Files
QYZH.InteractiveMagazine/QYZH.InteractiveMagazine.WorkService/Consumers/RabbitMQHostedService.cs
glz 16ad80f2a8 refactor(work-service): 重构自动铺码消费者,新增微信登录逻辑与配置更新
1. 更新微信公众号配置AppId和AppSecret
2. 修复微信授权控制器临时返回逻辑,启用真实登录流程
3. 新增期刊打印DTO类,定义铺码数据结构
4. 重构AutoDotCodeConsumer,新增依赖注入与完整铺码业务逻辑
5. 调整队列消费者接口与实现,添加取消令牌参数
6. 新增铺码回调响应实体类,完善异常处理与事务管理
2026-06-23 16:53:18 +08:00

128 lines
5.3 KiB
C#
Raw 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, ea.Body.ToArray(), 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, 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;
}
}
}