DEV Community

Cover image for MassTransit on RabbitMQ: what ends up in _error and _skipped, and how to move it back
Viktor
Viktor

Posted on Originally published at warrenops.io

MassTransit on RabbitMQ: what ends up in _error and _skipped, and how to move it back

Originally published at warrenops.io.

A consumer throws, MassTransit retries, and after the last attempt the message is gone from the endpoint. It is not lost: it sits in a queue next to it, shipping-label_error, with the exception in its headers. RabbitMQ's dead-lettering played no part in that, which is why the usual DLQ advice (x-death, dead-letter exchanges) does not quite fit. Here is what MassTransit does, what the headers say, and how to put the messages back once the bug is fixed.

The short version as a video (1:17, captions, no sound needed):

The topology MassTransit declares

For a receive endpoint shipping-label on RabbitMQ, MassTransit declares a queue shipping-label and an exchange of the same name bound to it. Message types get exchanges of their own (Shop.Contracts:CreateShippingLabel), bound to the endpoint's exchange when a consumer on that endpoint handles the type. Two more pairs appear the first time they are needed:

  • shipping-label_error: messages whose consumer threw, after all retries.
  • shipping-label_skipped: messages that arrived at the endpoint but that no consumer there handles, typically after a consumer was removed or renamed while publishers kept sending.
builder.Services.AddMassTransit(x =>
{
    x.AddConsumer<CreateShippingLabelConsumer>();
    x.UsingRabbitMq((context, cfg) =>
    {
        cfg.Host("rabbitmq", "/", h =>
        {
            h.Username("shipping");
            h.Password(builder.Configuration["RabbitMq:Password"]);
        });
        cfg.ReceiveEndpoint("shipping-label", e =>
        {
            e.ConfigureConsumer<CreateShippingLabelConsumer>(context);
        });
    });
});
Enter fullscreen mode Exit fullscreen mode

MassTransit publishes the message itself into the error queue, through the shipping-label_error exchange, and acknowledges the original. The broker never dead-lettered anything, so there is no x-death header and no dead-letter exchange on the queue. Tools and scripts that look for x-death find nothing to work with.

What the headers tell you

The error queue keeps MassTransit's JSON envelope as the body (application/vnd.masstransit+json: messageId, conversationId, sourceAddress, destinationAddress, messageType and your message) and adds headers that answer most of the postmortem:

Header Says
MT-Reason fault for the error queue, skip for the skipped queue
MT-Fault-ExceptionType the exception, e.g. Shop.Shipping.CarrierUnavailableException
MT-Fault-Message its message
MT-Fault-StackTrace the stack trace
MT-Fault-ConsumerType, MT-Fault-MessageType which consumer failed on which message type
MT-Fault-RetryCount how many retries ran before it gave up
MT-Fault-Timestamp, MT-Host-* when, and on which machine and MassTransit version

Group the error queue by MT-Fault-ExceptionType and the incident usually explains itself: five CarrierUnavailableExceptions from one outage of the carrier, two InvalidOperationExceptions that are real bugs in the data. Those two groups want different treatment.

MassTransit also publishes a Fault<CreateShippingLabel> event when this happens. Subscribing to it is the cleanest way to alert from inside your code; watching the depth of *_error queues is the way to alert from outside it.

Retry, redelivery, and when a message reaches the error queue

cfg.ReceiveEndpoint("shipping-label", e =>
{
    // second-level: back to the broker, delivered again after 5, 15 and 30 minutes
    e.UseDelayedRedelivery(r => r.Intervals(TimeSpan.FromMinutes(5), TimeSpan.FromMinutes(15), TimeSpan.FromMinutes(30)));
    // first-level: in memory, a few quick attempts
    e.UseMessageRetry(r => r.Intervals(200, 1000, 3000));
    e.ConfigureConsumer<CreateShippingLabelConsumer>(context);
});
Enter fullscreen mode Exit fullscreen mode

UseMessageRetry retries inside the consumer, in memory, while the message stays unacknowledged; keep it to seconds. UseDelayedRedelivery hands the message back to the broker and has it delivered again later; it needs a scheduler (cfg.UseDelayedMessageScheduler() with RabbitMQ's delayed-message-exchange plugin) and counts the rounds in the header MT-Redelivery-Count. Only when both are used up does the message go to _error, with MT-Fault-RetryCount saying how far it got.

Two consequences for a replay later. A transient failure (a carrier down for twenty minutes) should rarely reach the error queue if the redelivery intervals outlast typical outages; if it does, the intervals were too short. And a message that still carries MT-Redelivery-Count from its first life starts the next one with part of its budget spent.

Moving the messages back

The bug is fixed, the carrier is back. The messages belong in shipping-label again. What to get right:

  1. Send them to the endpoint's exchange, shipping-label, routing key empty. It is bound to the queue, so the message arrives exactly where MassTransit would deliver it. Publishing to the message type's exchange instead would reach every endpoint that subscribes to the type, which is a second delivery for all of them.
  2. Leave the body alone. MassTransit dispatches on the envelope's messageType. Unwrap the JSON to "simplify" it and the consumer no longer recognises the message, which then lands in _skipped.
  3. Remove the fault and redelivery headers: MT-Reason, MT-Fault-* and MT-Redelivery-Count, so the message starts with a full retry budget and a later fault is not mixed up with this one.
  4. Confirm before you acknowledge, and go slowly if the consumer just recovered.
  5. Skipped messages are different: replaying them to the same endpoint skips them again. Deploy the consumer for their type first, or send them to the endpoint that has it.

The management UI's "Move messages" (shovel plugin) covers point 1 for a whole queue, and nothing else: no selection by exception, headers untouched, no record of what was moved. A small tool with RabbitMQ.Client 7 does all of it:

// RabbitMQ.Client 7: with confirmation tracking, BasicPublishAsync returns once the broker confirmed
var factory = new ConnectionFactory { Uri = new Uri(rabbitUri) };
await using var connection = await factory.CreateConnectionAsync("error-replay");
await using var channel = await connection.CreateChannelAsync(
    new CreateChannelOptions(publisherConfirmationsEnabled: true, publisherConfirmationTrackingEnabled: true));

const string errorQueue = "shipping-label_error";
var endpoint = errorQueue[..^"_error".Length];
var moved = 0;
while (moved < limit && await channel.BasicGetAsync(errorQueue, autoAck: false) is { } result)
{
    var props = new BasicProperties(result.BasicProperties);
    props.Headers = result.BasicProperties.Headers?
        .Where(h => h.Key != "MT-Reason" && h.Key != "MT-Redelivery-Count" && !h.Key.StartsWith("MT-Fault-"))
        .ToDictionary(h => h.Key, h => h.Value);

    // to the endpoint's exchange, body untouched: MassTransit dispatches on the envelope's messageType
    await channel.BasicPublishAsync(endpoint, routingKey: "", mandatory: true, basicProperties: props, body: result.Body);
    await channel.BasicAckAsync(result.DeliveryTag, multiple: false); // only after the confirm
    moved++;
}
Enter fullscreen mode Exit fullscreen mode

Messages you do not want to move (the two real bugs) need a filter before the publish: read MT-Fault-ExceptionType and leave those unacknowledged; they go back to the error queue when the channel closes. Add a delay between publishes for a large queue, and log what you moved where.

Checklist

  • Know your endpoints' _error and _skipped queues and watch their depth, or subscribe to Fault<T>.
  • First-level retry for seconds, delayed redelivery for outages, so only real failures reach _error.
  • Group the error queue by MT-Fault-ExceptionType before you move anything.
  • Replay to the endpoint's exchange, body untouched, MT-Fault-*, MT-Reason and MT-Redelivery-Count removed, confirmed before acknowledged.
  • Skipped messages need their consumer first.

Where Warren fits

Warren recognises MassTransit's _error and _skipped queues as dead-letter queues, groups them by MT-Fault-ExceptionType, and says in one sentence which consumer faulted after how many retries. A replay goes to the endpoint's exchange with the fault and redelivery headers removed, the envelope untouched, confirmed message by message and throttled if you want, and every replay is in the audit log. Free for one broker.

Try Warren in a minute (one compose file, demo broker with real dead letters included) ยท Documentation

Top comments (0)