skip to content

What does DeadLetterPublishingRecoverer do, and what does the target topic and message look like?

level: middleimportance: must knowfreq 65%

answer

  1. default DLT = <topic>.DLT, same partition
  2. needs KafkaTemplate
  3. adds DLT_* headers (exception, original topic/offset)
  4. custom resolver → partition -1 to decouple
  5. failIfSendResultIsError = don't lose on failed publish

basics

~20 s

It's a recoverer that, after retries are exhausted, republishes the failed record to a dead-letter topic (default: original topic name + ".DLT") using a KafkaTemplate, adding headers describing the failure so the message isn't lost.

solid answer

~40 s

DeadLetterPublishingRecoverer is the ConsumerRecordRecoverer you give to DefaultErrorHandler (or @RetryableTopic uses one internally). When retries are exhausted it publishes the original record to a dead-letter topic via a KafkaTemplate. By default the DLT name is `<original-topic>.DLT` and it targets the **same partition number**, so the DLT should have at least as many partitions — otherwise you must supply a destination resolver (BiFunction of record+exception → TopicPartition), e.g. returning partition -1 to let Kafka choose. It enriches the message with headers like `KafkaHeaders.DLT_EXCEPTION_FQCN`, `DLT_EXCEPTION_MESSAGE`, `DLT_EXCEPTION_STACKTRACE`, `DLT_ORIGINAL_TOPIC/PARTITION/OFFSET/TIMESTAMP`, preserving the key/value/headers so a human or a repair job can inspect and replay. Without it, DefaultErrorHandler's default recoverer just logs and skips — losing the record.

code

java · 19 lines
java
@Bean
public DeadLetterPublishingRecoverer recoverer(KafkaTemplate<Object, Object> template) {
    var recoverer = new DeadLetterPublishingRecoverer(template,
        // decouple from source partition count: let Kafka pick the partition
        (record, ex) -> new TopicPartition(record.topic() + ".DLT", -1));
    // If the DLT publish itself fails, fail recovery instead of losing the record
    recoverer.setFailIfSendResultIsError(true);
    recoverer.setWaitForSendResultTimeout(Duration.ofSeconds(5));
    return recoverer;
}

@KafkaListener(topics = "orders.DLT", groupId = "orders-dlt-inspectors")
public void onDeadLetter(ConsumerRecord<String, String> rec) {
    String origTopic = new String(rec.headers()
        .lastHeader(KafkaHeaders.DLT_ORIGINAL_TOPIC).value());
    String exMsg = new String(rec.headers()
        .lastHeader(KafkaHeaders.DLT_EXCEPTION_MESSAGE).value());
    log.warn("Dead letter from {} failed with: {}", origTopic, exMsg);
}

go deeper

for a junior

Know it sends failed records to a <topic>.DLT so they aren't lost.

for a middle

Explain the KafkaTemplate dependency, default naming/partitioning, and the DLT_* headers.

for a senior

Cover custom destination resolvers, deserialization-failure handling, and making publishes fail-safe.

for a principal

Design DLT topology, retention/alerting, replay tooling, and transactional (EOS) recovery guarantees.

**Role.** A `DeadLetterPublishingRecoverer` implements `ConsumerRecordRecoverer` — the callback DefaultErrorHandler invokes once the BackOff is exhausted (or immediately for a fatal exception). Its job is to move the poison record somewhere durable instead of dropping it, so downstream processing can continue while the bad record is preserved for inspection/replay. It is the recommended recoverer; without one the container just logs and commits past the record (silent data loss). **How it publishes.** It needs a `KafkaTemplate` (or a template resolver, to publish different value types with different serializers). It sends a new `ProducerRecord` carrying the original key, value, and headers to a **dead-letter topic**. **Default destination.** The default `DestinationResolver` sends to `<originalTopic>.DLT` and to the **same partition index** as the source record. Consequence: the DLT must have **at least as many partitions** as the source topic, or the send fails. Common fixes: (a) auto-create the DLT with matching partitions, or (b) provide your own resolver — `new DeadLetterPublishingRecoverer(template, (record, ex) -> new TopicPartition(record.topic() + ".DLT", -1))` — where partition `-1` lets Kafka's partitioner pick, removing the partition-count coupling. **Headers added.** It stamps standard `KafkaHeaders` so the DLT record is self-describing: - `DLT_EXCEPTION_FQCN` — exception class name - `DLT_EXCEPTION_CAUSE_FQCN` — cause class - `DLT_EXCEPTION_MESSAGE` — the message - `DLT_EXCEPTION_STACKTRACE` — full stack trace - `DLT_ORIGINAL_TOPIC`, `DLT_ORIGINAL_PARTITION`, `DLT_ORIGINAL_OFFSET`, `DLT_ORIGINAL_TIMESTAMP`, `DLT_ORIGINAL_TIMESTAMP_TYPE`, `DLT_ORIGINAL_CONSUMER_GROUP` These let a repair tool know exactly where the record came from and why it failed, and support replay back to the original topic. **Deserialization failures.** If the record couldn't even be deserialized, the container passes a special record whose value/key holds the raw bytes plus a `DeserializationException`. DeadLetterPublishingRecoverer knows how to publish the raw bytes to the DLT (so you don't lose the payload just because it couldn't be parsed). This is why the DLT often needs byte-array-tolerant serialization. **Customization hooks.** - `setHeadersFunction(...)` — add/override headers on the outgoing DLT record. - `setFailIfSendResultIsError(true)` and `setWaitForSendResultTimeout(...)` — make the recoverer synchronous and fail (so the container retries recovery) if the DLT publish itself fails, instead of silently losing the record. - `setRetainExceptionHeader`, `setStripPreviousExceptionHeaders` — control header accumulation across multiple DLT hops. - `setPartitionSelector` / custom destination resolver — route by exception type (e.g. transient → retry topic, fatal → parking DLT). **Transactions.** If the consumer is in a Kafka transaction (EOS), the DLT publish and the offset commit are in the same transaction, so recovery is atomic with consumption. **When to use.** Any time you can't afford to silently drop failed records and want an operational audit trail + replay path. It's the standard partner to DefaultErrorHandler and is also what `@RetryableTopic` uses under the hood to forward between retry topics and to the final DLT. **Gotchas.** (1) Same-partition default silently requires matching partition counts — a frequent first-deploy failure. (2) If the DLT send fails and you didn't set `failIfSendResultIsError`, the record is still lost. (3) DLT can grow unbounded — set retention and monitoring/alerts. (4) Publishing full stack traces as headers can bloat messages.

  • You deploy and DLT publishes fail with an invalid-partition error. What's the likely cause and fix?
    The default destination resolver targets the same partition number as the source record, and the .DLT topic has fewer partitions. Either create the DLT with at least as many partitions, or supply a resolver returning partition -1 so Kafka's partitioner chooses.
  • How does the recoverer preserve a payload that failed to deserialize?
    On a DeserializationException the container hands the recoverer the raw bytes plus the exception; DeadLetterPublishingRecoverer publishes those original bytes to the DLT, so the unparseable payload is retained rather than lost.

saying these in an interview costs you the question

  • Claiming the DLT topic is auto-created with the right partition count in all setups (depends on broker auto-create + it defaults to same partition index)
  • Thinking a failed DLT publish is always safe (only if failIfSendResultIsError is set does recovery fail loudly)
  • Believing headers include the record but not the failure metadata (it adds exception class, message, stacktrace, original topic/offset)

context