skip to content

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%

answer

  1. not all deser failures are poison pills
  2. corrupt bytes = permanent → DLQ now
  3. registry down / schema-not-yet = transient → bounded retry
  4. tune retryable vs non-retryable exception sets
  5. per-cause metrics: infra vs data-quality vs schema-gov

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.

solid answer

~50 s

Not every deserialization failure is a poison pill. The key design decision is **classification by cause**: a malformed/truncated payload or wrong wire format is *permanently* undeserializable — retrying never helps, so skip-to-DLQ immediately. But a `KafkaAvroDeserializer` failing because the **Schema Registry is temporarily unreachable** (network blip, 5xx), or because a schema isn't registered *yet*, is *transient/recoverable* — discarding it would be false data loss; you should retry with backoff and only quarantine after exhausting retries. So a robust handler inspects the exception type/cause (`RestClientException` / connection errors vs. `SerializationException` on the bytes themselves), applies a retry policy to recoverable ones, and routes truly-corrupt records straight to the DLQ. The default error-handler classification (deserialization errors are non-retryable) is right for corrupt bytes but wrong for registry-outage errors, so you customize the retryable/non-retryable sets. Add per-cause metrics and alerting: a spike in registry errors is an infra incident; a spike in corrupt records is an upstream producer regression.

go deeper

for a junior

Know that some failures are temporary (registry down) and shouldn't be thrown away like corrupt records.

for a middle

Distinguish corrupt-bytes (skip) from registry-unavailable (retry) and know retries must be bounded.

for a senior

Implement cause-based classification, tune retryable/non-retryable sets with backoff, and route per cause.

for a principal

Define the fail-open-on-infra / fail-closed-on-bad-data policy, per-cause observability, and replay-safe deterministic classification org-wide.

## The core insight "Deserialization failure" is an umbrella over causes with **opposite correct responses**: | Cause | Nature | Correct response | |---|---|---| | Truncated/garbled bytes, wrong format on topic | **Permanent** — bytes are bad | Skip-to-DLQ now; retry is useless | | Schema Registry unreachable (timeout/5xx/conn refused) | **Transient** — infra problem | Retry with backoff; do NOT discard | | Schema ID not registered yet (race with producer) | **Transient/recoverable** | Retry / brief block; quarantine only after timeout | | Incompatible reader schema vs writer schema | **Semi-permanent** — config/evolution bug | Quarantine + alert; needs human fix | Treating all of them as "poison pill → skip" causes **false data loss**: during a 30-second Schema Registry outage you'd silently DLQ thousands of perfectly good records. Treating all as "retry" causes the **wedged-partition** poison-pill loop on genuinely corrupt bytes. ## Classifying in practice Inspect the exception chain: - **Corrupt/format errors:** `org.apache.kafka.common.errors.SerializationException` whose cause is a parse/format failure on the bytes themselves (e.g. Avro `IOException`, bad magic byte). These are **non-retryable** → DLQ. - **Registry-availability errors:** the Confluent deserializer wraps `io.confluent.kafka.schemaregistry.client.rest.exceptions.RestClientException` or raw `IOException`/`SocketTimeoutException` from the HTTP client. A 5xx or connection failure is **retryable**; a 404 (schema/subject not found) may be retryable-with-timeout (producer race) or a real config error. - **Compatibility errors:** schema-resolution failures indicate an evolution/config bug — quarantine and alert; they won't fix on retry. In Spring Kafka, `DefaultErrorHandler` lets you tune `addRetryableExceptions(...)` / `addNotRetryableExceptions(...)`. By default deserialization exceptions are not retryable — correct for corrupt bytes, but you must add registry-availability exceptions to the **retryable** set with a backoff (`ExponentialBackOff`) so transient outages don't dump good data into the DLQ. In Kafka Streams, a custom `DeserializationExceptionHandler` can branch: return `FAIL` (or block/retry) for registry-unavailable so the thread pauses and recovers when the registry returns, and route truly-corrupt records to a DLQ and return `CONTINUE`. ## Avoiding the wedge while retrying Retrying transient errors must be **bounded** — infinite retry on a misclassified permanent error recreates the poison-pill stall. Use exponential backoff with a max attempts/time cap, then fall through to DLQ. This bounds the blast radius: transient errors recover; persistent ones eventually quarantine. ## Observability and feedback loops Emit **separate metrics per cause**: - Registry-error rate → infra alert (page on-call; it affects *all* consumers). - Corrupt-record rate → data-quality alert (likely an upstream producer deployed a format change). - Compatibility-error rate → schema-governance alert (someone evolved a schema incompatibly). This turns one opaque "deserialization failed" log into actionable, routed signals. ## Replay and re-drive For recoverable causes, the records were never bad — once the registry recovers or the schema is registered, they deserialize fine, so no DLQ entry exists to re-drive. For corrupt/compatibility records in the DLQ, fixing the producer or registering a corrected schema lets a re-drive job replay the preserved original bytes back through the pipeline. Because the source log is immutable, classification must stay **deterministic** so a future replay re-encounters and re-classifies each record identically. ## Principal-level framing The architecture you're defending: **fail-open on transient infra, fail-closed on bad data** — never discard recoverable records, never loop forever on corrupt ones, and make every skip observable so silent loss is impossible to miss.

  • Why is treating a Schema Registry outage as a poison pill dangerous?
    The records are perfectly valid — they fail only because the registry is unreachable. Skipping them to a DLQ causes false data loss at scale during the outage. Registry-availability errors should be retried with backoff, not discarded.
  • How do you retry transient errors without recreating the wedged-partition stall?
    Bound the retries with exponential backoff and a max attempts/time cap. Transient causes recover within the budget; anything still failing falls through to the DLQ, so a misclassified permanent error can't loop forever.
  • In Spring Kafka, which knob lets you make registry-availability errors retryable while keeping corrupt-byte errors non-retryable?
    DefaultErrorHandler's addRetryableExceptions / addNotRetryableExceptions (with a configured BackOff). Deserialization errors are non-retryable by default; you add the registry/IO exceptions to the retryable set.

saying these in an interview costs you the question

  • Treating every deserialization failure identically as a poison pill to skip (causes false data loss on transient registry outages).
  • Retrying genuinely corrupt bytes (recreates the infinite poison-pill loop).
  • Unbounded retry on transient errors (wedges the partition just like a poison pill).
  • One generic 'deserialization failed' metric instead of per-cause signals (can't tell an infra outage from a producer regression).
  • Nondeterministic classification that diverges on replay.

context