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¶
- Python 3.11+,
pip install aio-pika; measured with RabbitMQ 4.3. The patterns apply to any broker. - A consumer, from processing RabbitMQ messages with aio-pika.
- Idempotent handlers, from making background jobs idempotent.
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.
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.
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.
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-limitfor 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.
Related¶
- Message Brokers & Event Streams — up to the topic overview.
- Tuning RabbitMQ prefetch for async consumers — why prefetch decides whether one poison message stalls the queue.
- Network I/O & Protocol Handling — the section overview.