From 1dd130693708fda4d86d9b2a22068001a3197f07 Mon Sep 17 00:00:00 2001 From: MilePrivate Date: Mon, 10 Aug 2026 13:37:10 +0200 Subject: [PATCH] Add RabbitMQ retry and dead-letter handling --- .../RabbitMqIntegrationEventPublisher.cs | 28 ++++++-- .../Services/Messaging/RabbitMqOptions.cs | 12 ++++ OrderProcessing.Api/appsettings.json | 8 ++- .../Configuration/RabbitMqOptions.cs | 12 ++++ .../RabbitMqEmailConsumer.cs | 68 +++++++++++++++---- OrderProcessing.EmailWorker/appsettings.json | 6 ++ 6 files changed, 115 insertions(+), 19 deletions(-) diff --git a/OrderProcessing.Api/Services/Messaging/RabbitMqIntegrationEventPublisher.cs b/OrderProcessing.Api/Services/Messaging/RabbitMqIntegrationEventPublisher.cs index b716fdf..f640c31 100644 --- a/OrderProcessing.Api/Services/Messaging/RabbitMqIntegrationEventPublisher.cs +++ b/OrderProcessing.Api/Services/Messaging/RabbitMqIntegrationEventPublisher.cs @@ -136,12 +136,30 @@ await channel.ExchangeDeclareAsync( autoDelete: false, cancellationToken: cancellationToken); + var emailQueueArguments = new Dictionary + { + ["x-queue-type"] = "quorum", + + ["x-delivery-limit"] = _options.DeliveryLimit, + + ["x-delayed-retry-type"] = "failed", + + ["x-delayed-retry-min"] = _options.RetryMinDelayMilliseconds, + + ["x-delayed-retry-max"] = _options.RetryMaxDelayMilliseconds, + + ["x-dead-letter-exchange"] = _options.DeadLetterExchangeName, + + ["x-dead-letter-routing-key"] = _options.DeadLetterRoutingKey + }; + await channel.QueueDeclareAsync( - queue: _options.EmailQueueName, - durable: true, - exclusive: false, - autoDelete: false, - cancellationToken: cancellationToken); + queue: _options.EmailQueueName, + durable: true, + exclusive: false, + autoDelete: false, + arguments: emailQueueArguments, + cancellationToken: cancellationToken); await channel.QueueBindAsync( queue: _options.EmailQueueName, diff --git a/OrderProcessing.Api/Services/Messaging/RabbitMqOptions.cs b/OrderProcessing.Api/Services/Messaging/RabbitMqOptions.cs index ccb95c6..65b596e 100644 --- a/OrderProcessing.Api/Services/Messaging/RabbitMqOptions.cs +++ b/OrderProcessing.Api/Services/Messaging/RabbitMqOptions.cs @@ -23,4 +23,16 @@ public sealed class RabbitMqOptions public string ClientProvidedName { get; set; } = "order-processing-api-publisher"; public int NetworkRecoveryIntervalSeconds { get; set; } = 5; + + public string DeadLetterExchangeName { get; set; } = "order-processing.dead-letter"; + + public string DeadLetterQueueName { get; set; } = "order-processing.email.dead-letter"; + + public string DeadLetterRoutingKey { get; set; } = "email.failed"; + + public int DeliveryLimit { get; set; } = 5; + + public int RetryMinDelayMilliseconds { get; set; } = 5000; + + public int RetryMaxDelayMilliseconds { get; set; } = 30000; } \ No newline at end of file diff --git a/OrderProcessing.Api/appsettings.json b/OrderProcessing.Api/appsettings.json index 4ebe231..faa3dd8 100644 --- a/OrderProcessing.Api/appsettings.json +++ b/OrderProcessing.Api/appsettings.json @@ -44,7 +44,7 @@ "PollingIntervalSeconds": 5 }, "RabbitMq": { - "Enabled": false, + "Enabled": true, "HostName": "localhost", "Port": 5672, "UserName": "guest", @@ -52,6 +52,12 @@ "VirtualHost": "/", "ExchangeName": "order-processing.events", "EmailQueueName": "order-processing.email", + "DeadLetterExchangeName": "order-processing.dead-letter", + "DeadLetterQueueName": "order-processing.email.dead-letter", + "DeadLetterRoutingKey": "email.failed", + "DeliveryLimit": 5, + "RetryMinDelayMilliseconds": 5000, + "RetryMaxDelayMilliseconds": 30000, "ClientProvidedName": "order-processing-api-publisher", "NetworkRecoveryIntervalSeconds": 5 } diff --git a/OrderProcessing.EmailWorker/Configuration/RabbitMqOptions.cs b/OrderProcessing.EmailWorker/Configuration/RabbitMqOptions.cs index 7cfad54..3fddb5f 100644 --- a/OrderProcessing.EmailWorker/Configuration/RabbitMqOptions.cs +++ b/OrderProcessing.EmailWorker/Configuration/RabbitMqOptions.cs @@ -23,4 +23,16 @@ public sealed class RabbitMqOptions public int NetworkRecoveryIntervalSeconds { get; set; } = 5; public ushort PrefetchCount { get; set; } = 1; + + public string DeadLetterExchangeName { get; set; } = "order-processing.dead-letter"; + + public string DeadLetterQueueName { get; set; } = "order-processing.email.dead-letter"; + + public string DeadLetterRoutingKey { get; set; } = "email.failed"; + + public int DeliveryLimit { get; set; } = 5; + + public int RetryMinDelayMilliseconds { get; set; } = 5000; + + public int RetryMaxDelayMilliseconds { get; set; } = 30000; } \ No newline at end of file diff --git a/OrderProcessing.EmailWorker/RabbitMqEmailConsumer.cs b/OrderProcessing.EmailWorker/RabbitMqEmailConsumer.cs index e37d4f9..2a8199d 100644 --- a/OrderProcessing.EmailWorker/RabbitMqEmailConsumer.cs +++ b/OrderProcessing.EmailWorker/RabbitMqEmailConsumer.cs @@ -189,12 +189,11 @@ await RejectPermanentFailureAsync( _logger.LogError( exception, "Temporary failure processing RabbitMQ " + - "message {MessageId}. The message will be requeued", + "message {MessageId}. The message will be retried", eventArgs.BasicProperties.MessageId); - await channel.BasicNackAsync( + await channel.BasicRejectAsync( deliveryTag: eventArgs.DeliveryTag, - multiple: false, requeue: true); } } @@ -204,20 +203,62 @@ private async Task RejectPermanentFailureAsync( BasicDeliverEventArgs eventArgs, Exception exception) { - _logger.LogError( - exception, - "Permanent failure processing RabbitMQ " + - "message {MessageId}. The message will not be requeued", - eventArgs.BasicProperties.MessageId); - - await channel.BasicNackAsync( - deliveryTag: eventArgs.DeliveryTag, - multiple: false, - requeue: false); + _logger.LogError(exception, + "Permanent failure processing RabbitMQ " + + "message {MessageId}. " + + "The message will be dead-lettered", + eventArgs.BasicProperties.MessageId); + + await channel.BasicRejectAsync( + deliveryTag: eventArgs.DeliveryTag, + requeue: false); } private async Task DeclareTopologyAsync(IChannel channel, CancellationToken cancellationToken) { + var emailQueueArguments = new Dictionary + { + ["x-queue-type"] = "quorum", + + ["x-delivery-limit"] = _options.DeliveryLimit, + + ["x-delayed-retry-type"] = "failed", + + ["x-delayed-retry-min"] = _options.RetryMinDelayMilliseconds, + + ["x-delayed-retry-max"] = _options.RetryMaxDelayMilliseconds, + + ["x-dead-letter-exchange"] = _options.DeadLetterExchangeName, + + ["x-dead-letter-routing-key"] = _options.DeadLetterRoutingKey + }; + + await channel.ExchangeDeclareAsync( + exchange: _options.DeadLetterExchangeName, + type: ExchangeType.Direct, + durable: true, + autoDelete: false, + cancellationToken: cancellationToken); + + var deadLetterQueueArguments = new Dictionary + { + ["x-queue-type"] = "quorum" + }; + + await channel.QueueDeclareAsync( + queue: _options.DeadLetterQueueName, + durable: true, + exclusive: false, + autoDelete: false, + arguments: deadLetterQueueArguments, + cancellationToken: cancellationToken); + + await channel.QueueBindAsync( + queue: _options.DeadLetterQueueName, + exchange: _options.DeadLetterExchangeName, + routingKey: _options.DeadLetterRoutingKey, + cancellationToken: cancellationToken); + await channel.ExchangeDeclareAsync( exchange: _options.ExchangeName, type: ExchangeType.Topic, @@ -230,6 +271,7 @@ await channel.QueueDeclareAsync( durable: true, exclusive: false, autoDelete: false, + arguments: emailQueueArguments, cancellationToken: cancellationToken); await channel.QueueBindAsync( diff --git a/OrderProcessing.EmailWorker/appsettings.json b/OrderProcessing.EmailWorker/appsettings.json index 6e3ee7c..4c8cdd7 100644 --- a/OrderProcessing.EmailWorker/appsettings.json +++ b/OrderProcessing.EmailWorker/appsettings.json @@ -13,6 +13,12 @@ "VirtualHost": "/", "ExchangeName": "order-processing.events", "EmailQueueName": "order-processing.email", + "DeadLetterExchangeName": "order-processing.dead-letter", + "DeadLetterQueueName": "order-processing.email.dead-letter", + "DeadLetterRoutingKey": "email.failed", + "DeliveryLimit": 5, + "RetryMinDelayMilliseconds": 5000, + "RetryMaxDelayMilliseconds": 30000, "ClientProvidedName": "order-processing-email-worker", "NetworkRecoveryIntervalSeconds": 5, "PrefetchCount": 1