feat(bll): обновлен MqServiceV2, добавлена обработка ошибок при выполнение делегата. Обновлен image rabbitMq до версии 4.2.0
This commit is contained in:
@@ -75,12 +75,26 @@ namespace PARR.BLL.Services.Implementations
|
|||||||
|
|
||||||
consumer.ReceivedAsync += async (ch, ea) =>
|
consumer.ReceivedAsync += async (ch, ea) =>
|
||||||
{
|
{
|
||||||
var content = Encoding.UTF8.GetString(ea.Body.ToArray());
|
//var content = Encoding.UTF8.GetString(ea.Body.ToArray());
|
||||||
|
var content = Encoding.UTF8.GetString(ea.Body.Span);
|
||||||
logger.LogDebug($"Получено сообщение: {content}");
|
logger.LogDebug($"Получено сообщение: {content}");
|
||||||
|
|
||||||
|
try
|
||||||
|
{
|
||||||
await messageHandler.Invoke(content);
|
await messageHandler.Invoke(content);
|
||||||
|
|
||||||
await consumerChannel.BasicAckAsync(ea.DeliveryTag, false);
|
await consumerChannel.BasicAckAsync(ea.DeliveryTag, false);
|
||||||
|
}
|
||||||
|
catch (Exception ex)
|
||||||
|
{
|
||||||
|
logger.LogError(ex, $"Ошибка при обработке сообщения: {content}");
|
||||||
|
|
||||||
|
//requeue: true - сообщение возвращаем в очередь в случае ошибки. Это немного опасно,
|
||||||
|
// если ошибка в формате сообщения, то она никогда не устранится и вечно будет добавляться/удаляться в очередь
|
||||||
|
//todo: тут можно реализовать DLQ с кол-вом попыток, после не успешных попыток, перемещать эти сообщения в другую очередь
|
||||||
|
// сделал false, выбрасывать на свежий воздух куда подальше эти кривые сообщения которые обработал с ошибками
|
||||||
|
await consumerChannel.BasicNackAsync(ea.DeliveryTag, false, requeue: false);
|
||||||
|
//await consumerChannel.BasicNackAsync(ea.DeliveryTag, false, requeue: true);
|
||||||
|
}
|
||||||
};
|
};
|
||||||
|
|
||||||
var consumerTag = await consumerChannel.BasicConsumeAsync(mqSettings.QueueName, false, consumer);
|
var consumerTag = await consumerChannel.BasicConsumeAsync(mqSettings.QueueName, false, consumer);
|
||||||
@@ -124,6 +138,12 @@ namespace PARR.BLL.Services.Implementations
|
|||||||
//props.Expiration = "60000";
|
//props.Expiration = "60000";
|
||||||
props.ContentType = "text/plain";//"application/json";
|
props.ContentType = "text/plain";//"application/json";
|
||||||
|
|
||||||
|
channel.BasicReturnAsync += async (sender, ea) =>
|
||||||
|
{
|
||||||
|
var body = Encoding.UTF8.GetString(ea.Body.Span);
|
||||||
|
logger.LogWarning($"Сообщение не было доставлено в очередь: {ea.Exchange} -> {ea.RoutingKey}. Тело: {body}");
|
||||||
|
};
|
||||||
|
|
||||||
foreach (var msg in msgList)
|
foreach (var msg in msgList)
|
||||||
{
|
{
|
||||||
var body = Encoding.UTF8.GetBytes(msg);
|
var body = Encoding.UTF8.GetBytes(msg);
|
||||||
@@ -157,11 +177,13 @@ namespace PARR.BLL.Services.Implementations
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
private async Task QueueDeclareAsync(IChannel channel, IMqSettings mqSettings)
|
private async Task QueueDeclareAsync(IChannel channel, IMqSettings mqSettings)
|
||||||
{
|
{
|
||||||
//https://www.rabbitmq.com/lazy-queues.html
|
//https://www.rabbitmq.com/lazy-queues.html
|
||||||
// lazy - ленивой очереди больше нет
|
// lazy - ленивой очереди больше нет
|
||||||
var args = new Dictionary<string, object?> { { "x-queue-mode", "lazy" } };
|
//var args = new Dictionary<string, object?> { { "x-queue-mode", "lazy" } };
|
||||||
|
var args = new Dictionary<string, object?> { };
|
||||||
|
|
||||||
//Объявляем очередь с которой будем работать.
|
//Объявляем очередь с которой будем работать.
|
||||||
//Если такой очереди ещё нет, то создатся.
|
//Если такой очереди ещё нет, то создатся.
|
||||||
|
|||||||
@@ -16,18 +16,13 @@ namespace PARR.GeneratorTemplatesWorker
|
|||||||
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
|
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
|
||||||
{
|
{
|
||||||
await generatorTemplate.StartAsync();
|
await generatorTemplate.StartAsync();
|
||||||
//while (!stoppingToken.IsCancellationRequested)
|
await Task.Delay(Timeout.Infinite, stoppingToken);
|
||||||
//{
|
|
||||||
// _logger.LogInformation("Worker running at: {time}", DateTimeOffset.Now);
|
|
||||||
// await Task.Delay(1000, stoppingToken);
|
|
||||||
//}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
public override Task StopAsync(CancellationToken cancellationToken)
|
public override async Task StopAsync(CancellationToken cancellationToken)
|
||||||
{
|
{
|
||||||
generatorTemplate.StopAsync().Wait();
|
await generatorTemplate.StopAsync();
|
||||||
|
await base.StopAsync(cancellationToken);
|
||||||
return base.StopAsync(cancellationToken);
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -27,6 +27,7 @@ services:
|
|||||||
|
|
||||||
parr-rabbitmq:
|
parr-rabbitmq:
|
||||||
image: rabbitmq:3.12.4-management
|
image: rabbitmq:3.12.4-management
|
||||||
|
#image: rabbitmq:4.2.0-management
|
||||||
#container_name: rabbitmq
|
#container_name: rabbitmq
|
||||||
hostname: parr-rabbitmq
|
hostname: parr-rabbitmq
|
||||||
environment:
|
environment:
|
||||||
|
|||||||
Reference in New Issue
Block a user