Files
QYZH.InteractiveMagazine/QYZH.InteractiveMagazine.WorkService/Consumers/RabbitMQHostedService.cs
glz 4142b621c3 feat: 新增工作服务项目并调整部分接口路由
1. 新增QYZH.InteractiveMagazine.WorkService工作服务项目,包含队列消费者框架、Hangfire定时任务支持
2. 修复PetService中未定义FeedPetOutput变量的问题
3. 调整JournalController的路由前缀从api/icr改为api
4. 将工作服务项目添加到解决方案中
2026-06-11 18:01:07 +08:00

83 lines
3.0 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.QueueName, stoppingToken)).ToList();
await Task.WhenAll(tasks);
}
private async Task StartConsumerAsync(string queueName, CancellationToken stoppingToken)
{
_logger.LogInformation("正在启动消费者: {QueueName}", queueName);
try
{
var rabbitMQService = _serviceProvider.GetRequiredService<IRabbitMQService>();
await rabbitMQService.ReceiveAsync(queueName, async (channel, ea) =>
{
// 每条消息创建独立 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);
await channel.BasicNackAsync(ea.DeliveryTag, false, true, 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);
}
}
}