skip to content

Deserialization Failure and Poison-Pill Handling

Surviving a record that will not deserialize: error-handling deserializers, Streams exception handlers, quarantine topics, and skipping forward. Interviewers ask because one bad message otherwise blocks a partition indefinitely.

part ofApache Kafkaoverview, primer and where to startread it →
on this pageshow

questions

5

What is a "poison pill" record in a Kafka consumer, and why can it stall a consumer group?

level: juniorimportance: must knowfreq 70%

answer

  1. deserialization runs inside poll()
  2. offset never advances → infinite retry
  3. blocks the whole partition in order
  4. can't try/catch your loop body
  5. fix at the deserializer layer

basics

~20 s

A poison pill is a record the consumer cannot deserialize (corrupt or wrong-format bytes). Deserialization happens inside poll(), so it throws every time, the offset never advances, and the consumer is stuck retrying the same record forever.

solid answer

~50 s

A "poison pill" is a message whose bytes the consumer's deserializer cannot turn into an object — e.g. a malformed Avro payload, a wrong schema ID, or plain text on a topic expecting JSON. In Kafka the deserializer runs inside KafkaConsumer.poll(), before your application code ever sees the record. So the failure throws a SerializationException out of poll() itself. Because the consumer hasn't processed (and therefore hasn't committed) the record, on the next poll it re-fetches the same offset and fails again. The result is an infinite retry loop that blocks all later records on that partition — the whole group makes no progress on that partition. You can't catch it in normal record-handling code because the exception happens before your loop body. The fix is to handle the failure at the deserialization layer (e.g. Spring's ErrorHandlingDeserializer) so the bad record is skipped and the offset advances.

go deeper

for a junior

Know the definition (a record you can't deserialize) and that it makes the consumer stuck retrying.

for a middle

Explain that deserialization happens inside poll(), so the offset can't advance and the partition wedges.

for a senior

Tie it to the exact exception (SerializationException out of poll), per-partition blocking, and the deserialization-layer fix.

for a principal

Reason about replay safety, key vs value pills, and choosing skip vs DLQ vs fail per data-quality SLA.

## The term A **poison pill** in Kafka is a record that a consumer **cannot deserialize**: the raw bytes on the topic can't be converted into the Java/Kotlin object the consumer expects. Common causes: - A producer wrote a different format than the consumer expects (e.g. plain `String` bytes on a topic the consumer reads as Avro). - A corrupt or truncated payload. - An Avro/Protobuf message whose **schema ID** (the bytes after the magic byte in the Confluent wire format) isn't resolvable, or whose schema is incompatible with the reader schema. - A `null`/empty payload where the deserializer doesn't tolerate it. ## Why it stalls the consumer — the mechanism In the Kafka client, **deserialization is not part of your application code** — it runs *inside* `KafkaConsumer.poll()`. The flow per poll: 1. The consumer fetches raw `byte[]` records from the broker for the assigned partitions. 2. For each record it calls the configured `key.deserializer` / `value.deserializer` to produce the typed `ConsumerRecord`. 3. Only *after* that does `poll()` return records to your loop. If step 2 throws (a `org.apache.kafka.common.errors.SerializationException`), the exception propagates **out of `poll()` itself**. Your `for (record in records)` body never runs for that batch. Crucially, the consumer's **position/offset is not advanced** for the failing record, and nothing is committed. On the **next** `poll()`, the consumer re-fetches starting at the same offset, deserialization fails again, and you get an **infinite loop**. Because Kafka delivers a partition's records in order, every record *after* the poison pill on that partition is also blocked — the partition (and effectively the group's progress on it) is wedged. ## Why you can't just try/catch your loop The naive instinct — wrap `record.value()` access in try/catch — doesn't help, because the throw happens **before** you ever get the record. The exception is thrown by `poll()`, not by reading the value. So the failure must be handled **at the deserialization layer**, where you can decide to (a) substitute a sentinel/null and let the offset advance, (b) route the bad bytes to a dead-letter queue, or (c) fail loudly and stop. ## The standard fix Spring Kafka's **`ErrorHandlingDeserializer`** wraps your real deserializer; when the delegate throws, it captures the exception in record headers and returns `null` instead of propagating, so `poll()` succeeds and the offset can advance. A `DefaultErrorHandler` / `DeadLetterPublishingRecoverer` then quarantines the bad record. In Kafka Streams, the analogous knob is `default.deserialization.exception.handler` (`LogAndContinueExceptionHandler` to skip, `LogAndFailExceptionHandler` to stop). ## Edge cases - A poison pill on **one partition** only blocks that partition; other partitions in the same consumer keep flowing. - Replaying from an earlier offset (e.g. after resetting offsets) will hit the same poison pill again unless you've handled it — handling must be **deterministic and replay-safe**. - The poison pill can be in the **key** deserializer too, not just the value.

  • Why doesn't wrapping your record-processing loop in try/catch fix a poison pill?
    Because the SerializationException is thrown inside poll() during deserialization, before any record is returned to your loop. Your catch block never runs; the throw escapes poll() itself, so it must be handled at the deserialization layer.
  • Does a poison pill on partition 0 stop the consumer from reading partition 3?
    No. Ordering and offset progress are per-partition, so partition 3 keeps flowing. Only partition 0 is wedged because its next offset can't advance past the bad record.

saying these in an interview costs you the question

  • Saying you can catch it in the consumer's record loop with try/catch (the throw happens inside poll, before the loop).
  • Claiming a poison pill blocks the entire topic across all partitions (it blocks only the partition it sits on).
  • Thinking auto-commit alone advances past it (the offset can't advance because poll itself fails).

context

open as a page

How does Spring Kafka's ErrorHandlingDeserializer work, and how does it cooperate with a dead-letter queue?

level: middleimportance: must knowfreq 60%

basics

~20 s

ErrorHandlingDeserializer 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.

open as a page

In Kafka Streams, how do you handle deserialization failures, and when would you choose LogAndContinue vs LogAndFail?

level: middleimportance: should knowfreq 50%

basics

~10 s

Set default.deserialization.exception.handler. LogAndContinueExceptionHandler logs the bad record and skips it (offset advances); LogAndFailExceptionHandler logs and stops the stream thread. Choose Continue for best-effort tolerance, Fail when no record may be silently dropped.

open as a page

After quarantining poison pills to a DLQ, how do you keep offset progress and replay/reprocessing safe and idempotent?

level: seniorimportance: should knowfreq 35%

basics

~20 s

Only commit the offset after the DLQ publish succeeds, so a crash mid-quarantine re-tries instead of skipping silently. Make handling deterministic per record and DLQ writes idempotent/keyed, so replaying the same bad record from an earlier offset produces the same quarantine, not duplicates or new gaps.

open as a page

How would you design a poison-pill strategy that distinguishes truly corrupt records from transient/recoverable deserialization failures, and routes each correctly?

level: principalimportance: should knowfreq 25%

basics

~20 s

Classify the failure: genuinely corrupt or wrong-format bytes are permanent — skip to DLQ immediately. Failures from a transiently unavailable Schema Registry or an unregistered-but-fixable schema are recoverable — retry with backoff and do NOT discard, because retrying corrupt bytes wastes effort and discarding recoverable ones causes false data loss.

open as a page