How would you design a poison-pill strategy that distinguishes truly corrupt records from transient/recoverable deserialization failures, and routes each correctly?
answer
- not all deser failures are poison pills
- corrupt bytes = permanent → DLQ now
- registry down / schema-not-yet = transient → bounded retry
- tune retryable vs non-retryable exception sets
- per-cause metrics: infra vs data-quality vs schema-gov
basics
~20 sClassify 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 sNot 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
Know that some failures are temporary (registry down) and shouldn't be thrown away like corrupt records.
Distinguish corrupt-bytes (skip) from registry-unavailable (retry) and know retries must be bounded.
Implement cause-based classification, tune retryable/non-retryable sets with backoff, and route per cause.
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.