using Microsoft.Extensions.Logging;
using RabbitMQ.Client;
using RabbitMQ.Client.Events;
using System.Text;
namespace QYZH.InteractiveMagazine.Infrastructure.MessageQueue;
///
/// RabbitMQ消息消费者基类
///
public abstract class RabbitMQConsumer : IDisposable
{
private readonly IConnection _connection;
private readonly ILogger _logger;
private IChannel? _channel;
private AsyncEventingBasicConsumer? _consumer;
///
/// 构造函数
///
/// RabbitMQ连接
/// 日志记录器
protected RabbitMQConsumer(IConnection connection, ILogger logger)
{
_connection = connection;
_logger = logger;
}
///
/// 启动消费
///
/// 队列名称
/// 消息处理委托
public async Task StartConsume(string queueName, Func handleMessage)
{
_channel = await _connection.CreateChannelAsync();
_consumer = new AsyncEventingBasicConsumer(_channel);
_consumer.ReceivedAsync += async (model, ea) =>
{
try
{
var body = ea.Body.ToArray();
var message = Encoding.UTF8.GetString(body);
await handleMessage(message);
await _channel.BasicAckAsync(ea.DeliveryTag, false);
_logger.LogInformation("消息消费成功 | 队列: {QueueName} | 消息: {Message}", queueName, message);
}
catch (Exception ex)
{
_logger.LogError(ex, "消息消费失败 | 队列: {QueueName}", queueName);
await _channel.BasicNackAsync(ea.DeliveryTag, false, true);
}
};
await _channel.BasicConsumeAsync(queue: queueName, autoAck: false, consumer: _consumer);
_logger.LogInformation("开始消费消息 | 队列: {QueueName}", queueName);
}
///
/// 释放资源
///
public void Dispose()
{
_channel?.DisposeAsync().GetAwaiter().GetResult();
GC.SuppressFinalize(this);
}
}