Files
QYZH.InteractiveMagazine/QYZH.InteractiveMagazine.WorkService/Consumers/RabbitMQHostedService.cs
glz d9933e4537 refactor: 重构RabbitMQ消费者体系,新增自动铺码功能
1. 新增AutoDotCodeConsumer自动铺码消费者,实现统一的交换机路由绑定
2. 重构RabbitMQ接收方法,支持交换机、路由键配置,添加死信队列处理
3. 重构现有JournalTaskReceiveConsumer,使用标准交换机路由配置
4. 更新JournalPageService,将打印逻辑替换为自动铺码消息发送
5. 调整枚举类型,重构任务类型命名
6. 更新实体和DTO,新增Prompt相关字段
7. 添加PrintToolV2.7配套工具和文档
8. 清理冗余的PrintJournalPageAsync接口
2026-06-22 16:56:05 +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());
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;
}
}
}