refactor: 拆分自动铺码服务为独立PrintWorker项目
1. 将原WorkService中的AutoDotCodeConsumer、RabbitMQHostedService、IQueueConsumer迁移至新的PrintWorker项目 2. 移除WorkService中冗余的PrintTool依赖文件复制配置 3. 将PrintToolV2.7相关资源移动到PrintWorker项目中统一管理 4. 从WorkService中移除AutoDotCodeConsumer的服务注册 5. 新增PrintWorker项目的完整宿主程序与配置文件
This commit is contained in:
@ -0,0 +1,105 @@
|
||||
using QYZH.InteractiveMagazine.Infrastructure.RabbitMQ;
|
||||
|
||||
namespace QYZH.InteractiveMagazine.PrintWorker.Consumers;
|
||||
|
||||
/// <summary>
|
||||
/// RabbitMQ 消费者后台服务。
|
||||
/// </summary>
|
||||
public class RabbitMQHostedService(IServiceProvider serviceProvider, ILogger<RabbitMQHostedService> logger) : BackgroundService
|
||||
{
|
||||
/// <summary>
|
||||
/// 启动所有已注册队列消费者。
|
||||
/// </summary>
|
||||
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;
|
||||
}
|
||||
|
||||
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) =>
|
||||
{
|
||||
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);
|
||||
|
||||
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);
|
||||
|
||||
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);
|
||||
}
|
||||
}
|
||||
|
||||
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;
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user