Skip to content

Handling Poison Messages in Async Consumers

A poison message is one that fails every time it is processed: malformed JSON, a reference to a deleted record, a payload that triggers a bug. The reflex response to a processing error — reject it back to the queue so it can be retried — turns a poison message into an infinite loop. Tested on RabbitMQ 4.3 with aio-pika, one malformed message among 100 in a classic queue, handled with nack(requeue=True): the 99 good messages were processed, and the bad one was redelivered 14,161 times in 3 seconds, burning a CPU core and flooding the logs. A quorum queue with x-delivery-limit: 3 did not stop it — 694 redeliveries in 2 s and nothing dead-lettered — because RabbitMQ 4 counts consumer crashes toward the limit but not explicit requeues; when the consumer crashed instead, the message was dead-lettered after its fourth delivery with reason delivery_limit. Rejecting with requeue=False to a dead-letter exchange moved it aside on the first failure. This guide classifies failures, bounds retries, and parks poison messages where they can be inspected.

Prerequisites

1. Recognise the redelivery loop

The loop needs only a handler that requeues on any exception:

async def handle(message: aio_pika.IncomingMessage) -> None:
    try:
        event = json.loads(message.body)
        await process(event)
        await message.ack()
    except Exception:
        await message.nack(requeue=True)       # poison: back to the head of the queue, forever

Measured: 14,161 redeliveries of one malformed message in 3 s. With a prefetch of 10, the good messages still got through, which makes the loop easy to miss — throughput looks fine while one core spins. With a prefetch of 1, or with several poison messages, the queue stops making progress entirely. The symptom in metrics is a redelivery rate far above the publish rate, and a consumer at full CPU with nothing to show for it.

Verify: chart the broker's redelivery rate (redeliver in RabbitMQ's message stats); outside incidents it should be near zero.

Deliveries of one poison message, by handling strategy 4 horizontal bars comparing classic queue, nack(requeue=True) with the others. Deliveries of one poison message, by handling strategy classic queue, nack(requeue=True) 14,161 in 3 s quorum, delivery-limit 3, nack requeue 694 in 2 s, not limited quorum, delivery-limit 3, consumer crash 4, then dead-lettered reject(requeue=False) to DLX 1, then dead-lettered RabbitMQ 4.3, aio-pika; the limit counts unacknowledged returns from crashes, not explicit requeues. Only an explicit decision in the handler stops a requeue loop.

2. Classify the failure before deciding

Not every failure is poison. A database timeout will probably succeed on retry; a JSON decode error never will. Decide per exception type:

class Permanent(Exception):
    """The message itself is bad; retrying cannot help."""


PERMANENT = (json.JSONDecodeError, KeyError, ValueError, Permanent)
TRANSIENT = (asyncio.TimeoutError, ConnectionError, asyncpg.PostgresConnectionError)


async def handle(message: aio_pika.IncomingMessage) -> None:
    try:
        event = parse(message.body)                 # raises Permanent on schema errors
        await process(event)
    except PERMANENT:
        await message.reject(requeue=False)         # straight to the dead-letter exchange
        return
    except TRANSIENT:
        await retry_later(message)                  # bounded, see step 3
        return
    await message.ack()

Validation and parsing errors are permanent by definition. Errors from dependencies are usually transient. Unknown exceptions are the hard case: treat them as transient but bounded, so a bug eventually lands the message in the dead-letter queue rather than looping. Tested, reject(requeue=False) on a queue with a dead-letter exchange moved the message aside immediately, with headers x-first-death-reason: rejected and x-first-death-queue naming the source.

Verify: a malformed test message ends up in the dead-letter queue after exactly one delivery.

3. Count attempts yourself

Because explicit requeues are not counted by the broker, count attempts in the message. Republish a copy with an incremented header and acknowledge the original; once the count reaches the limit, dead-letter it:

MAX_ATTEMPTS = 5


async def retry_later(message: aio_pika.IncomingMessage) -> None:
    attempt = int((message.headers or {}).get("x-attempt", 0)) + 1
    if attempt >= MAX_ATTEMPTS:
        await message.reject(requeue=False)                       # give up: dead-letter
        return
    delay_ms = min(1000 * 2 ** attempt, 60_000)
    await channel.default_exchange.publish(
        aio_pika.Message(
            message.body,
            headers={**(message.headers or {}), "x-attempt": attempt},
            expiration=delay_ms / 1000,                           # waits in the retry queue
            delivery_mode=aio_pika.DeliveryMode.PERSISTENT,
        ),
        routing_key="work.retry",       # no consumers; dead-letters back to "work" on expiry
    )
    await message.ack()

The retry queue has no consumers and is declared with x-dead-letter-exchange pointing back at the work queue, so a message waits there for its TTL and then returns — a delay with no timers in your process. Publish before acknowledging: if the process dies between the two, the message is duplicated, not lost, which idempotent handlers absorb. Quorum queues' x-delivery-limit remains useful as a backstop for the case it does cover: a message that crashes the consumer process every time.

Verify: a message that fails transiently four times and then succeeds is processed once; one that always fails reaches the dead-letter queue after MAX_ATTEMPTS deliveries.

Routing a failed message A flow of 5 stages. Routing a failed message handler fails classify the error permanent reject -> DLQ now transient republish x-attempt+1 delay queue TTL back to work queue attempt = max reject -> DLQ Every failure ends either in success or in the dead-letter queue, never in a loop.

4. Make the dead-letter queue useful

A dead-letter queue nobody looks at is a slow way to lose data. Give it what an operator needs, and a way back:

async def declare(channel: aio_pika.Channel) -> None:
    dlx = await channel.declare_exchange("work.dlx", aio_pika.ExchangeType.FANOUT, durable=True)
    dlq = await channel.declare_queue("work.dlq", durable=True)
    await dlq.bind(dlx)
    await channel.declare_queue("work", durable=True, arguments={
        "x-queue-type": "quorum",
        "x-dead-letter-exchange": "work.dlx",
        "x-delivery-limit": 10,                 # backstop for consumers that crash on a message
    })


async def replay_dlq(channel, limit: int = 100) -> int:
    dlq = await channel.get_queue("work.dlq")
    moved = 0
    while moved < limit and (msg := await dlq.get(fail=False)):
        await channel.default_exchange.publish(
            aio_pika.Message(msg.body, headers={"x-attempt": 0, "x-replayed": True}),
            routing_key="work",
        )
        await msg.ack()
        moved += 1
    return moved

Log the error and the message identifier when dead-lettering, so the DLQ entry can be matched to a stack trace. Alert on DLQ depth, not just on its existence. After fixing the bug, replay the dead-lettered messages in bounded batches with the attempt counter reset — the replay tool is as much a part of the design as the queue.

Verify: a test run produces a dead-lettered message, an alert fires, and the replay tool returns it to the work queue successfully after the fix.

5. Apply the same rules to Kafka and streams

Log-based brokers such as Kafka have no per-message requeue: a consumer that does not advance past a failing record blocks its whole partition. The equivalent of a dead-letter exchange is a dead-letter topic:

async def consume(consumer, producer) -> None:
    async for record in consumer:
        try:
            await process(parse(record.value))
        except PERMANENT as exc:
            await producer.send_and_wait("orders.dlt", record.value, key=record.key,
                                         headers=[("error", repr(exc).encode())])
        except TRANSIENT:
            await retry_with_backoff(lambda: process(parse(record.value)), attempts=5,
                                     on_give_up=lambda: producer.send_and_wait("orders.dlt", record.value))
        await consumer.commit({TopicPartition(record.topic, record.partition): record.offset + 1})

The offset is committed only after the record was processed or parked, so nothing is skipped silently; the commit rules are in committing Kafka offsets safely in async consumers. In-place retries block the partition while they back off, so keep their total time short and move longer retries to a retry topic consumed with a delay.

Verify: a poison record on one partition is moved to the dead-letter topic and the partition's lag recovers within the retry budget.

What should happen to a message that failed? A decision on Why did it fail with 4 outcomes. What should happen to a message that failed? Why did it fail? bad payload, validation dead-letter now retries cannot help dependency timeout, connection retry with backoff x-attempt header unknown exception retry, bounded then dead-letter consumer process crashed broker delivery-limit backstop Bounded retries plus a dead-letter queue turn poison into a ticket, not an outage.

Verification

Poison messages are handled when:

  • No handler requeues unconditionally.
  • Failures are classified, with permanent ones dead-lettered on first delivery.
  • Transient retries are counted and bounded in the message itself.
  • The dead-letter queue is monitored and replayable.

Diagnostic Hook: compare redelivery rate with publish rate per queue, and track dead-letter queue depth. A redelivery rate many times the publish rate is a requeue loop in progress; a dead-letter queue that grows steadily after a deploy is a new bug classifying good messages as bad.

Pitfalls & edge cases

  • nack(requeue=True) on every error. Measured: 14,161 redeliveries of one message in 3 s.
  • Relying on quorum x-delivery-limit for requeues. RabbitMQ 4 counts crashes, not explicit requeues.
  • Acknowledging before republishing a retry. A crash between them loses the message.
  • An unmonitored dead-letter queue. Messages are parked and forgotten.

Frequently Asked Questions

What is a poison message?

A message that fails processing every time, such as malformed JSON or a reference to data that no longer exists. Requeuing it creates an infinite redelivery loop.

Does RabbitMQ x-delivery-limit stop nack requeue loops?

Not on RabbitMQ 4.3 in testing: explicit nack with requeue was redelivered 694 times in 2 s with a limit of 3. The limit applied when the consumer crashed without acknowledging, dead-lettering the message after its fourth delivery.

How do I limit retries for a RabbitMQ message in aio-pika?

Republish a copy with an incremented attempt header, usually through a delay queue with a TTL, acknowledge the original, and reject with requeue=False once the count reaches the limit.

How do I handle poison messages in Kafka consumers?

Send the record to a dead-letter topic with the error in its headers, then commit the offset so the partition can move on. Keep in-place retries short, because they block the partition.