Files
QYZH.InteractiveMagazine/QYZH.InteractiveMagazine.Infrastructure/MessageQueue/RabbitMQConsumer.cs

73 lines
2.2 KiB
C#
Raw Normal View History

2026-06-01 13:42:40 +08:00
using Microsoft.Extensions.Logging;
using RabbitMQ.Client;
using RabbitMQ.Client.Events;
using System.Text;
namespace QYZH.InteractiveMagazine.Infrastructure.MessageQueue;
/// <summary>
/// RabbitMQ消息消费者基类
/// </summary>
public abstract class RabbitMQConsumer : IDisposable
{
private readonly IConnection _connection;
private readonly ILogger<RabbitMQConsumer> _logger;
private IChannel? _channel;
private AsyncEventingBasicConsumer? _consumer;
/// <summary>
/// 构造函数
/// </summary>
/// <param name="connection">RabbitMQ连接</param>
/// <param name="logger">日志记录器</param>
protected RabbitMQConsumer(IConnection connection, ILogger<RabbitMQConsumer> logger)
{
_connection = connection;
_logger = logger;
}
/// <summary>
/// 启动消费
/// </summary>
/// <param name="queueName">队列名称</param>
/// <param name="handleMessage">消息处理委托</param>
public async Task StartConsume(string queueName, Func<string, Task> 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);
}
/// <summary>
/// 释放资源
/// </summary>
public void Dispose()
{
_channel?.DisposeAsync().GetAwaiter().GetResult();
GC.SuppressFinalize(this);
}
}