A non-deserializable 'poison-pill' record keeps crashing your consumer in an infinite loop. Why does this happen with naive handling, and how do you fix it?
answer
- Deser happens inside poll(), before listener
- Bad record loops same offset forever
- ErrorHandlingDeserializer wraps delegate
- Failure → null payload + header exception
- DeserializationException is fatal → DLT immediately
basics
~20 sA poison pill is a record that can't be deserialized, so it throws before your listener even runs. The consumer keeps re-reading the same offset and looping forever. Fix it with ErrorHandlingDeserializer, which catches the failure and lets the error handler skip the record to a DLT.
solid answer
~40 sDeserialization happens inside poll(), before your @KafkaListener is invoked, so a malformed record throws a DeserializationException that the listener-level error handler can't recover normally — the consumer never advances past that offset and loops infinitely. The fix is ErrorHandlingDeserializer: wrap your real key/value deserializers with it (spring.deserializer.value.deserializer.delegate.class). On failure it doesn't throw out of poll(); instead it returns a null payload plus a DeserializationException in the record headers. Then DefaultErrorHandler sees the deserialization error, treats it as fatal/non-retryable, and immediately routes the record to the recoverer — typically a DeadLetterPublishingRecoverer that sends it to a DLT. This 'skips the poison pill' so the partition makes progress while preserving the bad record for inspection.
code
properties · 4 linesspring.kafka.consumer.value-deserializer=org.springframework.kafka.support.serializer.ErrorHandlingDeserializer
spring.deserializer.value.delegate.class=org.springframework.kafka.support.serializer.JsonDeserializer
spring.kafka.consumer.key-deserializer=org.springframework.kafka.support.serializer.ErrorHandlingDeserializer
spring.deserializer.key.delegate.class=org.apache.kafka.common.serialization.StringDeserializergo deeper
Know that a poison pill is an undeserializable record that loops the consumer, and ErrorHandlingDeserializer + a DLT fixes it.
Explain that deserialization happens in poll() before the listener, configure ErrorHandlingDeserializer delegates, and wire the DLT recoverer.
Discuss fatal classification, raw-bytes DLT publishing, failedDeserializationFunction, and key vs value failures.
Mandate ErrorHandlingDeserializer org-wide, design DLT inspection/replay and schema-governance to reduce poison pills at the source.
## What a poison pill is A **poison pill** is a record that cannot be processed no matter how many times you retry — classically one whose bytes can't be **deserialized** into your expected type (wrong schema, corrupt JSON, a record produced with an incompatible serializer). ## Why naive handling loops forever Key timing detail: **deserialization happens inside the consumer's `poll()` call**, *before* Spring ever invokes your `@KafkaListener`. The standard `KafkaConsumer` deserializes keys/values as it returns records. If the value deserializer throws, the exception propagates **out of `poll()`** — your listener code and its error handler never get a clean shot at the record. The consumer position hasn't advanced, so the next `poll()` returns the **same** bad record, which throws again: an **infinite crash loop** that also blocks every record behind it on that partition. ## The fix: `ErrorHandlingDeserializer` Spring's **`ErrorHandlingDeserializer`** wraps your actual deserializer (the *delegate*). Configuration (consumer props): ``` spring.kafka.consumer.key-deserializer: org.springframework.kafka.support.serializer.ErrorHandlingDeserializer spring.kafka.consumer.value-deserializer: org.springframework.kafka.support.serializer.ErrorHandlingDeserializer spring.deserializer.key.delegate.class: org.apache.kafka.common.serialization.StringDeserializer spring.deserializer.value.delegate.class: org.springframework.kafka.support.serializer.JsonDeserializer ``` When the delegate throws, `ErrorHandlingDeserializer` **swallows the exception** instead of letting it escape `poll()`. It returns a **`null` payload** and stashes the original bytes and a **`DeserializationException`** in special **record headers**. Now a normal record reaches the listener pipeline, so the container's error handler can act. ## How the record gets skipped to a DLT `DefaultErrorHandler` treats `DeserializationException` as a **fatal / not-retryable** exception by default (retrying a corrupt record is pointless). So it **skips retries** and invokes the recoverer immediately. Wire a **`DeadLetterPublishingRecoverer`** as that recoverer and the poison pill is **published to a DLT** with the failure metadata in headers. The consumer then commits past the offset — the partition is unblocked, and the bad record is preserved for inspection. ## Putting it together ``` @Bean DefaultErrorHandler errorHandler(KafkaTemplate<?, ?> template) { var recoverer = new DeadLetterPublishingRecoverer(template); return new DefaultErrorHandler(recoverer, new FixedBackOff(0L, 0L)); } ``` Deserialization errors are fatal, so they ignore the BackOff and go straight to the DLT regardless. ## Edge cases / gotchas - The `DeadLetterPublishingRecoverer` can send the **raw bytes** of the failed record (since the typed payload is null), which is exactly what you want for a deser failure. - You can register a `failedDeserializationFunction` on `ErrorHandlingDeserializer` to produce a placeholder object instead of null. - This also covers **key** deserialization failures, not just value. - If you *don't* use `ErrorHandlingDeserializer`, no listener-level error handler can save you from the poll-time loop — this is the single most common Kafka consumer footgun.
- Why can't an ordinary error handler catch a deserialization failure without ErrorHandlingDeserializer?Because deserialization throws inside poll(), before your listener and its error handler run. The exception escapes the poll loop, the offset never advances, and the consumer re-reads the same record forever. ErrorHandlingDeserializer moves the failure into the headers so the handler can act.
- Why is retrying a deserialization failure pointless, and how does DefaultErrorHandler reflect that?The bytes won't change between retries, so it will always fail. DefaultErrorHandler classifies DeserializationException as a fatal/not-retryable exception, skipping the backoff and going straight to the recoverer (DLT).
saying these in an interview costs you the question
- Saying a try/catch in the listener can fix it — the failure happens before the listener runs.
- Suggesting you retry the poison pill — retries can't fix corrupt bytes; it must be skipped/DLT'd.
- Forgetting to also wrap the key deserializer when keys can be malformed.
- Thinking the consumer auto-skips bad records — without ErrorHandlingDeserializer it loops forever.