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 MqServiceV2 : IMqService { private readonly ILogger logger; private IConnection? consumerConnection; private IChannel? consumerChannel; // реализация для RabbitMQ.Client 7.1.2 public MqServiceV2(ILogger logger) { this.logger = logger; } public async Task InitConsumerAsync(IMqSettings mqSettings, MqMessageHandlerDelegate messageHandler) { logger.LogInformation($"Устанавливаю соединение с RabbitMQ: {mqSettings.HostName}, {mqSettings.QueueName}"); var factory = new ConnectionFactory { HostName = mqSettings.HostName, UserName = mqSettings.User, Password = mqSettings.Password, AutomaticRecoveryEnabled = true }; try { consumerConnection = await factory.CreateConnectionAsync(); consumerChannel = await consumerConnection.CreateChannelAsync(); // кол-во сообщений которые можно обрабатывать за раз if (mqSettings.PrefetchCount.HasValue) { await consumerChannel.BasicQosAsync(0, mqSettings.PrefetchCount.Value, false); } // создаем очередь, вдруг ее еще нет await QueueDeclareAsync(consumerChannel, mqSettings); //await consumerChannel.QueueBindAsync(mqSettings.QueueName, "default", "routingKey"); consumerChannel.CallbackExceptionAsync += async (ch, ea) => { var _channel = (IChannel?)ch; // может тут если канал упал переоткрывать его logger.LogError(ea.Exception, $"Ошибка при обработке очереди MqService сообщений из очереди. Канал открыт?: {_channel?.IsOpen}"); //if (consumerChannel.IsClosed) //{ // //todo: этот момент проверить, успешно ли он переоткроет канал и очередь будет дальше обрабатываться // await consumerChannel.CloseAsync(); // consumerChannel = await consumerConnection.CreateChannelAsync(); // logger.LogInformation(ea.Exception, $"Переоткрыл канал RabbitMq."); //} }; var consumer = new AsyncEventingBasicConsumer(consumerChannel); //consumer.HandleChannelShutdownAsync = async (c, ea) =>{ }; consumer.ReceivedAsync += async (ch, ea) => { var content = Encoding.UTF8.GetString(ea.Body.ToArray()); logger.LogDebug($"Получено сообщение: {content}"); await messageHandler.Invoke(content); await consumerChannel.BasicAckAsync(ea.DeliveryTag, false); }; var consumerTag = await consumerChannel.BasicConsumeAsync(mqSettings.QueueName, false, consumer); return true; } catch (Exception ex) { logger.LogError(ex, $"Ошибка при получении сообщений из очереди: {mqSettings.HostName}, {mqSettings.QueueName}"); return false; } } public async Task SendAsync(IMqSettings mqSettings, string[] msgList) { var factory = new ConnectionFactory { HostName = mqSettings.HostName, UserName = mqSettings.User, Password = mqSettings.Password, AutomaticRecoveryEnabled = true }; // индекс текущей отправки в очередь из массива msgList // нужен для формирования списка неотправленных сообщений var sendIdx = 0; try { using (var connection = await factory.CreateConnectionAsync()) using (var channel = await connection.CreateChannelAsync()) { // создаем очередь, вдруг ее еще нет await QueueDeclareAsync(channel, mqSettings); var props = new BasicProperties(); //храним на диске props.DeliveryMode = DeliveryModes.Persistent; //время жизни, мс props.Expiration = "60000"; props.ContentType = "text/plain";//"application/json"; foreach (var msg in msgList) { var body = Encoding.UTF8.GetBytes(msg); await channel.BasicPublishAsync( exchange: "", routingKey: mqSettings.QueueName, //TODO: проверить mandatory: true, может указать его в false mandatory: true, 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() }; } } private async Task QueueDeclareAsync(IChannel channel, IMqSettings mqSettings) { //https://www.rabbitmq.com/lazy-queues.html // lazy - ленивой очереди больше нет var args = new Dictionary { { "x-queue-mode", "lazy" } }; //Объявляем очередь с которой будем работать. //Если такой очереди ещё нет, то создатся. //Если нет, нужно параметры типа durable должны совпадать иначе будет ошибка. //В целом эти параметры можно посмотреть в админке RabbitMQ await channel.QueueDeclareAsync( queue: mqSettings.QueueName, durable: true, exclusive: false, autoDelete: false, arguments: args ); } public async ValueTask DisposeAsync() { if (consumerChannel != null) { if (consumerChannel.IsOpen) { await consumerChannel.CloseAsync(); } await consumerChannel.DisposeAsync(); } if (consumerConnection != null) { if (consumerConnection.IsOpen) { await consumerConnection.CloseAsync(); } await consumerConnection.DisposeAsync(); } } } }