添加项目文件。
This commit is contained in:
@ -0,0 +1,72 @@
|
||||
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);
|
||||
}
|
||||
}
|
||||
@ -0,0 +1,47 @@
|
||||
using Microsoft.Extensions.Logging;
|
||||
using RabbitMQ.Client;
|
||||
using System.Text;
|
||||
|
||||
namespace QYZH.InteractiveMagazine.Infrastructure.MessageQueue;
|
||||
|
||||
/// <summary>
|
||||
/// RabbitMQ消息发布器
|
||||
/// </summary>
|
||||
public class RabbitMQPublisher
|
||||
{
|
||||
private readonly IConnection _connection;
|
||||
private readonly ILogger<RabbitMQPublisher> _logger;
|
||||
|
||||
/// <summary>
|
||||
/// 构造函数
|
||||
/// </summary>
|
||||
/// <param name="connection">RabbitMQ连接</param>
|
||||
/// <param name="logger">日志记录器</param>
|
||||
public RabbitMQPublisher(IConnection connection, ILogger<RabbitMQPublisher> logger)
|
||||
{
|
||||
_connection = connection;
|
||||
_logger = logger;
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// 发布消息
|
||||
/// </summary>
|
||||
/// <param name="exchange">交换机名称</param>
|
||||
/// <param name="routingKey">路由键</param>
|
||||
/// <param name="message">消息内容</param>
|
||||
public async Task PublishMessage(string exchange, string routingKey, string message)
|
||||
{
|
||||
await using var channel = await _connection.CreateChannelAsync();
|
||||
|
||||
var body = Encoding.UTF8.GetBytes(message);
|
||||
|
||||
await channel.BasicPublishAsync(
|
||||
exchange: exchange,
|
||||
routingKey: routingKey,
|
||||
body: body
|
||||
);
|
||||
|
||||
_logger.LogInformation("消息已发布 | 交换机: {Exchange} | 路由键: {RoutingKey} | 消息: {Message}",
|
||||
exchange, routingKey, message);
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user