DEV Community

Cover image for Poison Messages in Kafka: Retries, a Dead-Letter Queue, and the Replay Shortcut I'd Fix
Onkar Deokate
Onkar Deokate

Posted on

Poison Messages in Kafka: Retries, a Dead-Letter Queue, and the Replay Shortcut I'd Fix

TL;DR: One bad message can stop a Kafka consumer, because offsets are committed in order: if record #41 can never be processed, record #42 never gets processed either. In TxFlow I handle this with bounded retries plus exponential backoff, a dead-letter topic, and a DLQ handler that makes failures visible and replayable. This post covers the design, the order of operations that keeps it safe, and the replay shortcut I'd fix before calling it production-ready.


Continuing my TxFlow posts. Earlier: why I moved from SQS to Kafka and making consumers idempotent.

The problem: poison messages

A poison message is a record your consumer will never successfully process: malformed JSON, a missing field, a user that doesn't exist, a currency you don't support.

With SQS or Sidekiq, a failing message just sits in its own retry cycle while everything else keeps flowing. Kafka is different. A consumer group tracks one offset per partition, and committing offset 42 means "everything up to 42 is done". If you can't get past 41, you can't commit, and the partition is stuck. Consumer lag grows, and every user whose events hash to that partition waits.

You can reproduce this in TxFlow with one command:

echo '{bad json' | docker compose exec -T redpanda \
  rpk topic produce payments.initiated -X brokers=redpanda:9092 -k user_999
Enter fullscreen mode Exit fullscreen mode

Step 1: separate "try again" from "never going to work"

Not every failure deserves a retry:

  • Transient: database timeout, broker hiccup, downstream 503. Retrying often works.
  • Permanent: can't parse, fails validation, business rule rejects it. Retrying just wastes time while the partition waits.

A simplified sketch of the consumer framework's shape:

MAX_ATTEMPTS = 3

def process_with_retries(msg, handler):
    try:
        event = parse_and_validate(msg.value())
    except ValidationError as e:
        return send_to_dlq(msg, reason="invalid", error=e, attempts=0)  # no retries

    for attempt in range(1, MAX_ATTEMPTS + 1):
        try:
            return handler(event)
        except Exception as e:
            if attempt == MAX_ATTEMPTS:
                return send_to_dlq(msg, reason="exhausted", error=e, attempts=attempt)
            time.sleep(0.2 * 2 ** (attempt - 1))  # 0.2s, 0.4s...
Enter fullscreen mode Exit fullscreen mode

TxFlow retries up to 3 times with exponential backoff, inside the consumer. That's a deliberate trade-off. While a record is being retried, the rest of its partition waits. At demo volume, a second of blocking is fine. At high volume, you'd use retry topics instead (payments.retry.30s, payments.retry.5m, and so on), so the main topic keeps moving while failed records wait their turn somewhere else.

Step 2: a useful DLQ envelope

"Sent to DLQ" is only useful if someone can later work out what failed, where and why. The envelope should carry enough to debug and replay without digging through logs:

{
  "original_topic": "payments.initiated",
  "partition": 2,
  "offset": 41,
  "key": "user_999",
  "consumer_group": "wallet",
  "reason": "exhausted",
  "error": "InsufficientFunds: user_999",
  "attempts": 3,
  "failed_at": "2026-09-29T10:12:03Z",
  "payload": { "...": "original event, untouched" }
}
Enter fullscreen mode Exit fullscreen mode

consumer_group matters most. The same event can be fine for four groups and fail for one. When you replay it, you want only that one group to process it again.

Step 3: order of operations

The sequence that keeps this safe:

  1. Produce the envelope to payments.dlq
  2. Wait for the broker to acknowledge it (flush)
  3. Then commit the original offset

If you commit first and crash before the DLQ write lands, the failure disappears without a trace. If you write to the DLQ first and crash before committing, the worst case is a duplicate DLQ entry. That's easy to dedupe and much easier to live with than a lost failure.

TxFlow's DLQ topic has one partition. DLQ traffic is low, and a single ordered stream is easier to inspect and replay.

Step 4: make failures visible

A DLQ nobody looks at is just a slower way to lose data. TxFlow's dlq-handler service:

  • consumes payments.dlq,
  • stores each failure in a dead_letter_events table in Postgres,
  • exposes GET /dlq for the dashboard, which shows recent failures with search, filters and a Replay button,
  • sits next to consumer lag per group on the same dashboard, because a growing DLQ and a growing lag are the two signals that something's wrong.

Not everything needs a DLQ

The analytics consumer has no DLQ at all. It increments counters, and if an event can't be counted after retries, it isn't parked anywhere. A DLQ is an operational commitment: someone has to look at it. I only give one to consumers where a skipped event actually costs something (wallet, fraud, notify, audit).

The replay shortcut I'd fix

Today, the dashboard's Replay button sends the payment back through POST /payment. That's fine for a demo, but it's the wrong layer, and it's the first thing on my improvement list:

  • the API either rejects it (409, if the original idempotency key is reused) or produces a fresh event,
  • a fresh event gets a new event_id, so it doesn't match any consumer's dedup keys. All five groups process it again, not just the one that failed. A failure in the notify consumer could end up re-running the wallet debit,
  • the link to the original event is lost, which makes debugging harder.

The better design re-publishes the original event (same event_id, same payload) to a replay or retry topic that only the failed consumer group reads. Every other group's idempotency stays intact, and the audit trail shows it was a replay.

Checklist

If you're adding a DLQ to a Kafka consumer:

  • Send permanent failures (parsing, validation) straight to the DLQ, without retrying
  • Keep in-consumer retries few and short, or move to retry topics
  • Write to the DLQ, wait for the ack, then commit the offset
  • Put consumer_group, topic/partition/offset, error and attempt count in the envelope
  • Store DLQ entries somewhere queryable, with a dashboard or alert
  • Replay to the failed group only, with the original event ID

Question for you: do you use in-consumer retries or retry topics? And what's the strangest poison message you've found in a DLQ?

Top comments (0)