How does Spring Kafka's ErrorHandlingDeserializer work, and how does it cooperate with a dead-letter queue?
answer
- delegating/decorator deserializer
- returns null + stashes exception in headers
- poll() succeeds → container handles it
- DefaultErrorHandler → DeadLetterPublishingRecoverer
- DLT topic carries ORIGINAL raw bytes
basics
~20 sErrorHandlingDeserializer wraps your real deserializer. If the delegate throws, it catches the error, returns null, and stashes the exception and raw bytes in record headers. poll() then succeeds, so a DefaultErrorHandler/DeadLetterPublishingRecoverer can route the bad record to a DLQ and advance the offset.
solid answer
~40 s`ErrorHandlingDeserializer` is a delegating deserializer you configure via `ErrorHandlingDeserializer.VALUE_DESERIALIZER_CLASS` (and the KEY variant). It calls the real deserializer; if that throws, instead of letting the exception escape `poll()`, it returns a `null` value and writes the exception plus the original bytes into headers (`ErrorHandlingDeserializer.VALUE_DESERIALIZER_EXCEPTION_HEADER`, etc.). Because `poll()` now succeeds, the failure surfaces in normal listener-container error handling rather than as a wedged partition. The container's `DefaultErrorHandler` sees the deserialization exception is non-retryable and hands the record to a `DeadLetterPublishingRecoverer`, which republishes the **original raw bytes** (not the failed-to-deserialize object) to a `<topic>.DLT` topic, preserving headers so you can inspect the cause. Then the offset is committed and the consumer moves on. This converts a fatal, looping failure into a skip-and-quarantine.
go deeper
Know it's a wrapper that turns a deserialization crash into a skippable null so the consumer keeps going.
Explain the delegate config, the exception headers, and the handoff to DefaultErrorHandler + DeadLetterPublishingRecoverer.
Detail header preservation, original-bytes republish, ByteArraySerializer on the DLT, and null-vs-tombstone disambiguation.
Design the end-to-end quarantine: DLT topic strategy, triage headers, replay tooling, and non-retryable classification policy.
## The problem it solves A raw deserialization failure throws **out of `KafkaConsumer.poll()`**, so the offset can't advance and the partition wedges (the poison-pill problem). `ErrorHandlingDeserializer` moves that failure from *inside* `poll()` to a place where Spring's listener container can handle it. ## How it works — delegation You don't set your real deserializer directly. Instead: ```properties spring.kafka.consumer.value-deserializer=org.springframework.kafka.support.serializer.ErrorHandlingDeserializer spring.kafka.consumer.properties.spring.deserializer.value.delegate.class=io.confluent.kafka.serializers.KafkaAvroDeserializer ``` `ErrorHandlingDeserializer` is a **decorator**. On each record it calls `delegate.deserialize(...)`: - **Success:** returns the deserialized object as normal. - **Failure (delegate throws):** it does **not** rethrow. Instead it returns `null` as the value and attaches the caught exception (and the original `byte[]`) to the record via headers: - `ErrorHandlingDeserializer.VALUE_DESERIALIZER_EXCEPTION_HEADER` — a serialized `DeserializationException` carrying the cause and the original data. - (key variant: `KEY_DESERIALIZER_EXCEPTION_HEADER`.) Because it returns instead of throwing, `poll()` **succeeds** and returns the (null-valued, header-tagged) record to the container. ## Cooperating with the error handler and DLQ The listener container runs a `DefaultErrorHandler`. When a record carries a deserialization-exception header (or the listener throws because the value is null), the handler: 1. Recognizes deserialization failures as **non-retryable** (retrying can't fix bad bytes) — they're in the default "not retryable" set, so it doesn't loop. 2. Invokes the configured **recoverer**, typically `DeadLetterPublishingRecoverer`. `DeadLetterPublishingRecoverer`: - Reads the **original raw bytes** out of the header (critical — you can't re-serialize an object that never deserialized). - Publishes them to a dead-letter topic, by default `<originalTopic>.DLT`, same partition by default. - Copies headers, adding diagnostic headers (`kafka_dlt-exception-fqcn`, `kafka_dlt-exception-message`, `kafka_dlt-original-topic`, `kafka_dlt-original-offset`, etc.) so a human/operator can triage. After the recoverer succeeds, the container **commits the offset** and processing continues — the partition is unblocked. ## Edge cases and gotchas - **Null vs. legitimate tombstone:** since failures surface as `null` value, your listener must distinguish a real Kafka tombstone (`null` value, used for compaction deletes) from a deserialization failure. Check the exception header, not just `value == null`. - **DLQ also depends on a working serializer:** the DLQ producer must be able to publish the raw bytes — use `ByteArraySerializer` for the DLT producer so re-serialization can't fail again. - **Key vs value:** configure both `KEY_DESERIALIZER_CLASS` and `VALUE_DESERIALIZER_CLASS` delegates if either side can be malformed. - **Headers must survive:** if an intermediate system strips headers, the recoverer loses the original bytes. Keep the path header-preserving. - **Replay safety:** the behavior is deterministic — replaying the same bad record produces the same skip-to-DLT, so reprocessing from an earlier offset is safe.
- Why must the DLT carry the original raw bytes rather than the deserialized object?Because the object never deserialized — there is no object to re-serialize. The recoverer pulls the original byte[] from the exception header and republishes it (typically with a ByteArraySerializer) so the payload is preserved for triage.
- How does a listener tell a deserialization failure apart from a legitimate compaction tombstone, since both arrive as null?It inspects the deserialization-exception header (VALUE_DESERIALIZER_EXCEPTION_HEADER). A real tombstone has null value with no exception header; a poison pill has the header set.
saying these in an interview costs you the question
- Saying ErrorHandlingDeserializer retries the deserialization (it doesn't — bad bytes can't be fixed by retry; it skips/quarantines).
- Claiming the DLQ stores the deserialized object (it stores the original raw bytes).
- Forgetting that a deserialization failure surfaces as a null value, which can be confused with a tombstone.
- Setting your real deserializer directly instead of as the delegate behind ErrorHandlingDeserializer.