Handling Poison Messages in Kafka
Kafka has no built-in dead-letter queue, and its per-partition ordering means a record that always fails blocks every record behind it. This guide builds poison-message handling for Kafka consumers — retry topics, a dead-letter topic, and replay — as part of Dead-Letter Queues & Poison Messages in Queue Fundamentals & Architecture.
Problem Statement
A billing consumer reads invoice-events from a 24-partition topic. One morning a producer bug emits a record with a malformed currency code. The consumer throws, does not commit, re-polls, and throws again — forever. Partition 17 stops advancing; every invoice event behind the bad record (about 40,000 by the time anyone notices) waits. The consumer's error logs are flooded, but the lag alert only fires 50 minutes later because total group lag stays modest. You want the bad record moved aside within seconds, transient failures retried without blocking the partition, full context preserved for triage, and a safe way to replay fixed records.
Prerequisites
- Kafka consumer code with manual offset commits (
enable.auto.commit=false). - Permission to create topics: one or more retry topics and a dead-letter topic per consumer group.
- An understanding of whether the consumer needs strict per-key ordering (it changes the design; see Step 5).
- Idempotent processing, since retried records may be processed more than once.
Step 1 — Classify Failures Before Deciding Where Records Go
Three categories need three different treatments:
- Deserialization failures (bad bytes, schema mismatch): will never succeed. Dead-letter immediately.
- Permanent processing failures (validation errors, missing referenced entity that will never exist): dead-letter immediately.
- Transient failures (database timeout, downstream 503): retry after a delay.
class PermanentError(Exception): ...
class TransientError(Exception): ...
def classify(exc: Exception) -> str:
if isinstance(exc, (DeserializationError, SchemaError, ValidationError, PermanentError)):
return "dead-letter"
if isinstance(exc, (DatabaseTimeout, httpx.TimeoutException, TransientError)):
return "retry"
return "retry" # unknown: retry a bounded number of times, then dead-letter
The default matters: unknown exceptions go to retry with a bounded number of attempts, so a new kind of bug does not silently dead-letter valid records on the first blip — nor block the partition forever.
Step 2 — Create Retry and Dead-Letter Topics
Name topics after the consumer group, not the source topic alone, because different consumers of the same topic fail independently.
for t in billing.invoice-events.retry-1m billing.invoice-events.retry-10m; do
kafka-topics.sh --bootstrap-server $B --create --topic $t --partitions 24 --replication-factor 3 \
--config retention.ms=604800000 # 7 days
done
kafka-topics.sh --bootstrap-server $B --create --topic billing.invoice-events.dlt \
--partitions 6 --replication-factor 3 --config retention.ms=2592000000 # 30 days to triage
Match the retry topics' partition count to the main topic if you want records to keep their key-to-partition mapping; the DLT can be smaller since it carries little traffic.
Step 3 — Publish Failures with Context, Then Commit
When a record fails, produce it to the chosen topic with headers describing the failure, wait for the produce to be acknowledged, and only then commit the original offset. Committing first would lose the record if the produce failed.
def route_failure(producer, msg, exc: Exception, attempt: int) -> None:
decision = classify(exc)
if decision == "retry" and attempt < 3:
target = "billing.invoice-events.retry-1m" if attempt == 1 else "billing.invoice-events.retry-10m"
else:
target = "billing.invoice-events.dlt"
headers = [
("x-original-topic", msg.topic().encode()),
("x-original-partition", str(msg.partition()).encode()),
("x-original-offset", str(msg.offset()).encode()),
("x-attempt", str(attempt + 1).encode()),
("x-error-class", type(exc).__name__.encode()),
("x-error", str(exc)[:1000].encode()),
("x-failed-at", str(int(time.time() * 1000)).encode()),
("x-not-before", str(int((time.time() + delay_for(target)) * 1000)).encode()),
]
producer.produce(target, key=msg.key(), value=msg.value(), headers=(msg.headers() or []) + headers)
producer.flush(10) # must succeed before we commit
def handle(consumer, producer, msg):
attempt = int(dict(msg.headers() or []).get("x-attempt", b"1"))
try:
process(deserialize(msg.value()))
except Exception as exc:
route_failure(producer, msg, exc, attempt)
consumer.commit(message=msg, asynchronous=False)
Keeping the original key preserves the key's partition in the retry topic, and the original bytes (not the deserialized object) are what goes to the DLT, so deserialization bugs can be reproduced exactly.
Step 4 — Consume Retry Topics with Delays
A retry consumer must not process a record before its delay has passed. Rather than sleeping per record, pause the partition until the head record's x-not-before time, then resume.
def retry_loop(consumer, producer):
consumer.subscribe(["billing.invoice-events.retry-1m", "billing.invoice-events.retry-10m"])
while True:
msg = consumer.poll(1.0)
if msg is None or msg.error():
continue
not_before = int(dict(msg.headers()).get("x-not-before", b"0")) / 1000
wait = not_before - time.time()
if wait > 0:
tp = TopicPartition(msg.topic(), msg.partition(), msg.offset())
consumer.pause([tp]); consumer.seek(tp) # re-read this record later
schedule_resume(consumer, tp, wait) # resume after the delay
continue
handle(consumer, producer, msg)
Because records in a retry topic are written in time order, the head record is always the next one due; pausing on it does not delay records that are due earlier. Spring Kafka's non-blocking retries and similar libraries implement this pattern; the delay strategy mirrors exponential backoff with jitter.
Step 5 — Decide What Ordering You Can Give Up
Moving a failed record to a retry topic lets later records for the same key overtake it. For many consumers (idempotent upserts, independent events) that is fine. For per-key ordered processing it is not, and there are two options:
- Block the key, not the partition. Keep an in-memory (or small persistent) set of keys that currently have records in retry. Records for those keys go straight to the retry topic behind the failed one, preserving their relative order, while other keys continue.
- Block the partition deliberately for a bounded time, then dead-letter the head and alert — accepting that the whole partition pauses briefly.
blocked_keys: set[bytes] = load_blocked_keys() # keys with records currently in retry
def handle_ordered(consumer, producer, msg):
if msg.key() in blocked_keys:
route_to_retry_behind(producer, msg) # keep order for this key
else:
try:
process(deserialize(msg.value()))
except Exception as exc:
blocked_keys.add(msg.key())
route_failure(producer, msg, exc, attempt=1)
consumer.commit(message=msg, asynchronous=False)
The key is unblocked when its last retry record succeeds or is dead-lettered.
The ordering rules this preserves are discussed in Message Ordering Guarantees.
Step 6 — Triage and Replay from the Dead-Letter Topic
The DLT is only useful if someone looks at it and can replay what was fixed. Alert on any new DLT records, and build a replay tool that republishes selected records to the main topic with a marker.
# Inspect recent dead letters with their headers
kcat -b $B -t billing.invoice-events.dlt -C -o -20 -e \
-f 'offset=%o key=%k error=%h\n'
def replay(dlt_offsets: list[int], dry_run: bool = True) -> None:
for rec in read_dlt(dlt_offsets):
hdr = dict(rec.headers())
if dry_run:
print(rec.key(), hdr[b"x-error-class"], hdr[b"x-original-offset"])
continue
producer.produce(hdr[b"x-original-topic"].decode(), key=rec.key(), value=rec.value(),
headers=[("x-replayed-from-dlt", str(rec.offset()).encode())])
producer.flush()
Replay only after the fix is deployed, and rate-limit large replays. The replayed record is a new record at the end of the main topic, so if ordering matters, confirm no later records for the same key have been applied — the sequence-number approach in handling out-of-order events with sequence numbers makes this safe.
Verification
Produce a poison record into staging and confirm the partition keeps moving:
echo '{"invoice_id":"x","currency":"???"}' | kcat -b $B -t invoice-events -P -k inv-test -p 17
# Partition 17 lag should return to ~0 within seconds; the DLT should hold one record
kafka-consumer-groups.sh --bootstrap-server $B --describe --group billing | grep ' 17 '
kcat -b $B -t billing.invoice-events.dlt -C -o -1 -e -f '%k %h\n'
Gotchas & Edge Cases
Committing before producing. Commit only after the retry/DLT produce is acknowledged, or a crash between the two loses the record.
Retry topics that never drain. A downstream outage sends everything to retry; retry consumers then fail and forward everything to the DLT. Use a circuit breaker to pause the main consumer during outages instead — see circuit breakers for worker dependencies.
Schema registry deserialization. Records that fail Avro/Protobuf deserialization cannot be inspected with the consumer's own deserializer; keep raw bytes in the DLT and inspect with a byte-level tool.
Exactly-once pipelines. With Kafka transactions, produce to retry/DLT topics inside the same transaction as the offset commit.
FAQ
Why not just skip the bad record and log it? A log line is easy to miss and hard to replay. A DLT keeps the exact bytes, the error, and the origin, and makes replay a routine operation.
How quickly should a poison record be moved aside? Immediately for deserialization and validation failures — there is nothing to wait for. For unknown exceptions, after the first in-place retry fails; holding the partition for more than a few seconds turns one bad record into lag for every record behind it. The lag alert that finally fired after 50 minutes in the scenario should have been a per-partition "offset not advancing" alert firing within a minute or two.
What should the DLT alert include? The consumer group, error class, count in the window, and a link to a query or tool that lists the records with their headers. On-call needs to know within a minute whether this is one malformed record or a producer emitting nothing but bad data.
How many retry topics should I have? Two or three with increasing delays (seconds, minutes, tens of minutes) covers most transient failures. More adds complexity without much benefit.
Does Kafka Connect have this built in?
Yes, for connectors: errors.tolerance=all with errors.deadletterqueue.topic.name routes failed records to a DLT with context headers. Application consumers must implement it themselves or use a framework that does.
Related
- Dead-Letter Queues & Poison Messages — DLQ design across brokers.
- Alerting on Dead-Letter Queue Growth — making dead letters visible.
- Per-Entity Ordering with Kafka Partition Keys — the ordering this pattern must respect.
- Retrying Only Transient Errors by Exception Type — the classification in Step 1.