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 logger; private IConnection? connection; private IModel? channel; public MqService(ILogger 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 { { "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(); // формируем список неотправленных элементов 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 { { "x-queue-mode", "lazy" } }; channel.QueueDeclare( queue: mqSettings.QueueName, durable: true, exclusive: false, autoDelete: false, arguments: args ); logger.LogInformation($"Соединение с RabbitMQ установлено: {mqSettings.HostName}, {mqSettings.QueueName}"); channel.CallbackException += (ch, ea) => { var _channel = (IModel?)ch; // может тут если канал упал переоткрывать его logger.LogError(ea.Exception, $"Ошибка при обработке очереди MqService сообщений из очереди. Канал открыт: {_channel?.IsOpen}"); }; var consumer = new AsyncEventingBasicConsumer(channel); consumer.Received += async (ch, ea) => { var content = Encoding.UTF8.GetString(ea.Body.ToArray()); logger.LogInformation($"Получено сообщение: {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(); } } }