Topic 6.2
Retry Topics, Dead-Letter Topics, and Poison Messages
In one line
A poison message (one that always fails) must not block a partition forever. Route failing records to delayed retry topics (for example 5 s, 30 s, 5 min, 30 min), and after the last attempt to a dead-letter topic (DLT) with headers describing the original topic, partition, offset, exception and attempt count. Replay from the DLT after fixing the cause.
Think of it like this
A sorting office with a "problem parcels" shelf. A parcel with an unreadable address is set aside with a note explaining the problem, so the conveyor keeps moving. Parcels that just need another attempt later go to a "try again this afternoon" bin.
Key ideas
- 01
Retry topic chain:
orders.retry.5s,orders.retry.30s,orders.retry.5m,orders.retry.30m, thenorders.dlt. Each retry consumer waits until the record's due time (timestamp + delay) before processing, pausing the partition rather than sleeping inside poll. - 02
Headers to add: original topic, partition, offset and timestamp; exception class and message (and a stack trace hash); attempt number; first failure time; consumer group. They make triage and replay possible.
- 03
Spring Kafka provides this out of the box:
@RetryableTopic(non-blocking retry topics with backoff, DLT) orDefaultErrorHandlerwithDeadLetterPublishingRecoverer(blocking retries, then DLT). - 04
Ordering trade-off: once a record moves to a retry topic, later records with the same key continue on the main topic, so order for that key changes. If per-key order matters, either block the partition (in-place retries), or also park subsequent records for the same key (track "keys in retry" and route their new records to the retry path too).
- 05
DLT operations: alert on DLT growth, provide a tool to inspect, fix and replay records (republish to the main topic with the original key), and set DLT retention long enough for humans to act.
Code & diagrams
Spring Kafka non-blocking retries with exponential backoff and a DLT.
@RetryableTopic(
attempts = "5", // 1 original + 4 retries
backoff = @Backoff(delay = 5_000, multiplier = 6.0, maxDelay = 1_800_000),
exclude = { DeserializationException.class, ValidationException.class }, // straight to DLT
dltStrategy = DltStrategy.FAIL_ON_ERROR,
autoCreateTopics = "false")
@KafkaListener(topics = "commerce.orders", groupId = "payment-service")
public void onOrder(OrderCreated evt) {
paymentClient.authorize(evt); // throws on 503
}
@DltHandler
public void onDlt(OrderCreated evt, @Header(KafkaHeaders.EXCEPTION_MESSAGE) String error,
@Header(KafkaHeaders.ORIGINAL_OFFSET) long offset) {
log.error("order {} dead-lettered at offset {}: {}", evt.orderId(), offset, error);
metrics.counter("orders.dlt").increment();
}Interview problem
The problem
A poison message blocks everything
One malformed message makes every consumer attempt fail, the partition is stuck, and lag grows. Design retry topics, backoff, a DLT and a replay process so the consumer never gets stuck.
The interviewer follows up
Why not just skip bad messages and log them?
When it breaks
Deserialization errors crash the listener container
What you see
The consumer restarts, re-polls the same bad record, crashes again: an infinite loop, the partition stops, lag grows.
Fix & prevent
Use ErrorHandlingDeserializer (Spring) or catch deserialization in the client and route the raw bytes to the DLT.
Explain it without notes
What metadata should a DLT record carry and why?
Practice
Build the Spring retry-topic example and publish one record that always fails and one that fails twice then succeeds.
Trade-offs
- ↔
Non-blocking retries keep throughput but reorder; blocking retries keep order but stall partitions.
Done when you can
I can design retry topics and DLTs with metadata, alerts and replay.