skip to content

A bad record causes a deserialization exception that crashes your @KafkaListener in an infinite loop. How do you handle this poison pill?

level: seniorimportance: must knowfreq 65%

answer

  1. exception thrown inside poll(), before listener
  2. offset never commits → same record forever
  3. ErrorHandlingDeserializer wraps the real deserializer
  4. delegate.class = KafkaAvroDeserializer
  5. DefaultErrorHandler + DeadLetterPublishingRecorder → .DLT

basics

~20 s

Wrap your deserializer in Spring Kafka's ErrorHandlingDeserializer. It catches the exception during deserialization, hands a null payload plus the failure to your error handler, and lets you skip or DLT the record instead of looping forever.

solid answer

~40 s

A deserialization exception is thrown inside the consumer's poll loop, before your listener runs, so a normal try/catch in the listener can't see it; the container retries the same offset endlessly — a poison pill. The fix is Spring Kafka's ErrorHandlingDeserializer: set value.deserializer to ErrorHandlingDeserializer and put the real deserializer under spring.deserializer.value.function or spring.deserializer.value.delegate.class (KafkaAvroDeserializer). It catches the failure, passes a null value to the listener and stashes the exception/raw bytes in a header (DeserializationExceptionHeader). Pair it with a DefaultErrorHandler plus a DeadLetterPublishingRecorder so the bad record is routed to a <topic>.DLT and the offset advances. This separates infrastructure (can't decode) failures, which must skip, from business failures, which you may retry. Without ErrorHandlingDeserializer the container's SeekToCurrent/DefaultErrorHandler can't recover because the offset is never successfully consumed.

go deeper

for a junior

Recognize the symptom: a bad message makes the consumer loop forever.

for a middle

Name ErrorHandlingDeserializer and the delegate.class property as the fix.

for a senior

Wire ErrorHandlingDeserializer + DefaultErrorHandler + DLT and separate infra vs business failures.

for a principal

Design retry/DLT policy that tolerates transient registry outages without discarding valid data, plus the Streams equivalent.

## Why a poison pill is special The Kafka consumer deserializes records **inside `poll()`**, before any of your code runs. If `KafkaAvroDeserializer` throws (corrupt bytes, missing magic byte, unreachable registry for that ID, schema mismatch), the exception propagates out of the container's fetch. Crucially: - Your `@KafkaListener` body never executes, so a `try/catch` there is useless. - The offset is never committed because the record was never successfully processed. - On the next poll the container re-reads the **same** offset, fails again — an **infinite loop / poison pill** that stalls the whole partition. ## The fix: ErrorHandlingDeserializer Spring Kafka provides `org.springframework.kafka.support.serializer.ErrorHandlingDeserializer`. You set it as the actual `value.deserializer` and delegate to the real one: ``` spring.kafka.consumer.value-deserializer=org.springframework.kafka.support.serializer.ErrorHandlingDeserializer spring.kafka.properties.spring.deserializer.value.delegate.class=io.confluent.kafka.serializers.KafkaAvroDeserializer ``` Now when the delegate throws, `ErrorHandlingDeserializer` **catches** it and instead of failing the poll it: 1. Returns a `null` value to the listener container. 2. Stores the original exception and raw bytes in record headers (`ErrorHandlingDeserializer.VALUE_DESERIALIZER_EXCEPTION_HEADER`). Because deserialization now "succeeds" (with null), the container can apply its normal **error handling** path. ## Routing the bad record With a `DefaultErrorHandler` (the modern replacement for the old SeekToCurrentErrorHandler) plus a `DeadLetterPublishingRecorder`, a record that arrives null-due-to-deserialization is recognized and published to a **dead-letter topic** (`<topic>.DLT` by default), and the offset advances so the partition unblocks. You typically configure `FixedBackOff`/`ExponentialBackOff` with a finite number of retries before DLT. ## Separating failure classes This design enforces a key distinction: - **Infrastructure / deserialization failures** (can't decode the bytes at all) — retrying is pointless; skip or DLT immediately. - **Business failures** (decoded fine, but processing threw) — may be transient; retry with backoff. ErrorHandlingDeserializer cleanly captures the first class so your error handler can treat them as non-retryable. ## Edge cases & gotchas - You can also set `spring.deserializer.value.function` to a `Function<byte[]?, T>` to substitute a default value instead of null. - Do the same for the key (`spring.deserializer.key.delegate.class`) if keys are also Avro. - A transiently unreachable registry can masquerade as a poison pill; tune retries so you don't DLT records that would decode once the registry returns. Distinguishing transient registry outages from genuinely corrupt records is the subtle part. - In Kafka Streams the analog is `default.deserialization.exception.handler` (e.g. `LogAndContinueExceptionHandler` vs `LogAndFailExceptionHandler`).

  • Why can't a try/catch inside the @KafkaListener method handle a deserialization error?
    Deserialization happens during poll(), before the listener is invoked. The exception never reaches your method body, so the catch never fires; the container loops on the uncommitted offset instead.
  • How do you avoid DLT-ing a perfectly good record when the Schema Registry is briefly down?
    A registry outage throws during deserialization and looks like a poison pill. Use finite retries with backoff so transient registry failures recover, and only route to DLT after retries are exhausted — ideally distinguishing connectivity errors from genuine decode errors.
  • What's the Kafka Streams equivalent?
    Set default.deserialization.exception.handler to LogAndContinueExceptionHandler (skip) or LogAndFailExceptionHandler (stop), or a custom handler — Streams has no @KafkaListener container, so it uses this DSL-level hook.

saying these in an interview costs you the question

  • Claiming a try/catch in the listener fixes it — the exception fires before the listener runs.
  • Retrying deserialization failures indefinitely — corrupt bytes never decode; you just stall the partition.
  • Forgetting to advance/commit the offset, leaving the consumer group stuck on one partition.
  • Treating a transient registry outage as a permanent poison pill and DLT-ing valid records.

context