161 lines
6.3 KiB
C#
161 lines
6.3 KiB
C#
using Microsoft.Extensions.Logging;
|
||
using PARR.BLL.Contracts.Interfaces;
|
||
using PARR.BLL.Domain.Mq;
|
||
using PARR.BLL.Services.Interfaces;
|
||
using RabbitMQ.Client;
|
||
using RabbitMQ.Client.Events;
|
||
using System.Text;
|
||
|
||
namespace PARR.BLL.Services.Implementations
|
||
{
|
||
|
||
internal class MqService : IMqService
|
||
{
|
||
private readonly ILogger<MqService> logger;
|
||
|
||
private IConnection? connection;
|
||
private IModel? channel;
|
||
|
||
public MqService(ILogger<MqService> logger)
|
||
{
|
||
this.logger = logger;
|
||
}
|
||
|
||
public MqSendResult Send(IMqSettings mqSettings, string[] msgList)
|
||
{
|
||
var factory = new ConnectionFactory
|
||
{
|
||
HostName = mqSettings.HostName,
|
||
UserName = mqSettings.User,
|
||
Password = mqSettings.Password
|
||
};
|
||
|
||
// индекс текущей отправки в очередь из массива msgList
|
||
// нужен для формирования списка неотправленных сообщений
|
||
var sendIdx = 0;
|
||
|
||
try
|
||
{
|
||
using (var connection = factory.CreateConnection())
|
||
using (var channel = connection.CreateModel())
|
||
{
|
||
//https://www.rabbitmq.com/lazy-queues.html
|
||
//Рекомендовано при работе порциями использовать ленивые очереди. сообщения не используют память
|
||
//Формируем соответствующий аргумент
|
||
var args = new Dictionary<string, object> { { "x-queue-mode", "lazy" } };
|
||
|
||
//Объявляем очередь с которой будем работать.
|
||
//Если такой очереди ещё нет, то создатся.
|
||
//Если нет, нужно параметры типа durable должны совпадать иначе будет ошибка.
|
||
//В целом эти параметры можно посмотреть в админке RabbitMQ
|
||
|
||
channel.QueueDeclare(
|
||
queue: mqSettings.QueueName,
|
||
durable: true,
|
||
exclusive: false,
|
||
autoDelete: false,
|
||
arguments: args
|
||
);
|
||
|
||
var props = channel.CreateBasicProperties();
|
||
//храним на диске
|
||
props.DeliveryMode = 2;
|
||
//время жизни, мс
|
||
//props.Expiration = "60000";
|
||
|
||
foreach (var msg in msgList)
|
||
{
|
||
var body = Encoding.UTF8.GetBytes(msg);
|
||
channel.BasicPublish(exchange: "", routingKey: mqSettings.QueueName, basicProperties: props, body: body);
|
||
|
||
sendIdx++;
|
||
logger.LogDebug($"Отправлено сообщение в очередь: {mqSettings.HostName} {mqSettings.QueueName}, {msg}");
|
||
}
|
||
}
|
||
|
||
return new MqSendResult { IsSuccess = true };
|
||
}
|
||
catch (Exception ex)
|
||
{
|
||
var notSendMsg = new List<string>();
|
||
|
||
// формируем список неотправленных элементов
|
||
for (int i = sendIdx; i < msgList.Length; i++)
|
||
notSendMsg.Add(msgList[i]);
|
||
|
||
logger.LogError(ex, $"Ошибка при отправке сообщения в очередь {mqSettings.HostName} {mqSettings.QueueName}, {string.Join(',', notSendMsg)}");
|
||
|
||
return new MqSendResult { IsSuccess = false, NotSendMessages = notSendMsg.ToArray() };
|
||
}
|
||
}
|
||
|
||
public bool InitConsumer(IMqSettings mqSettings, MqMessageHandlerDelegate messageHandler)
|
||
{
|
||
logger.LogInformation($"Устанавливаю соединение с RabbitMQ: {mqSettings.HostName}, {mqSettings.QueueName}");
|
||
|
||
var facory = new ConnectionFactory
|
||
{
|
||
HostName = mqSettings.HostName,
|
||
UserName = mqSettings.User,
|
||
Password = mqSettings.Password,
|
||
AutomaticRecoveryEnabled = true,
|
||
DispatchConsumersAsync = true
|
||
};
|
||
|
||
try
|
||
{
|
||
//using var connection = facory.CreateConnection();
|
||
//using var channel = connection.CreateModel();
|
||
|
||
connection = facory.CreateConnection();
|
||
channel = connection.CreateModel();
|
||
|
||
var args = new Dictionary<string, object> { { "x-queue-mode", "lazy" } };
|
||
|
||
channel.QueueDeclare(
|
||
queue: mqSettings.QueueName,
|
||
durable: true,
|
||
exclusive: false,
|
||
autoDelete: false,
|
||
arguments: args
|
||
);
|
||
|
||
logger.LogInformation($"Соединение с RabbitMQ установлено: {mqSettings.HostName}, {mqSettings.QueueName}");
|
||
|
||
|
||
var consumer = new AsyncEventingBasicConsumer(channel);
|
||
consumer.Received += async (ch, ea) =>
|
||
{
|
||
var content = Encoding.UTF8.GetString(ea.Body.ToArray());
|
||
logger.LogDebug($"Получено сообщение: {content}");
|
||
|
||
await messageHandler.Invoke(content);
|
||
|
||
channel.BasicAck(ea.DeliveryTag, false);
|
||
await Task.Yield();
|
||
};
|
||
|
||
string consumerTag = channel.BasicConsume(mqSettings.QueueName, false, consumer);
|
||
|
||
return true;
|
||
|
||
}
|
||
catch (Exception ex)
|
||
{
|
||
logger.LogError(ex, $"Ошибка при получении сообщений из очереди: {mqSettings.HostName}, {mqSettings.QueueName}");
|
||
|
||
return false;
|
||
}
|
||
}
|
||
|
||
public void Dispose()
|
||
{
|
||
channel?.Close();
|
||
connection?.Close();
|
||
|
||
channel?.Dispose();
|
||
connection?.Dispose();
|
||
}
|
||
}
|
||
}
|