skip to content

A streaming ingestion consumer crashes on the same record after every restart — how do you isolate it?

level: seniorimportance: should knowfreq 48%

answer

  1. same position, same error, every time
  2. the stream stalls behind one message
  3. not every failure deserves diverting
  4. an outage would dead-letter everything
  5. diverting one record breaks that key's order

basics

~20 s

Classify the failure first. If it is deterministic, bound the attempts, divert that record with its key, position and error into a dead-letter destination, advance past it, and alert. Never divert on transient failures — an outage would dead-letter the whole stream.

solid answer

~50 s

The symptom — same position, same exception, every restart — is a poison record: a message that fails deterministically, so the consumer can never advance past it and the stream stalls behind it. The fix has three parts. **Classify** the error: a parse failure, a schema mismatch or a constraint violation is deterministic and will never succeed; a sink timeout or a connection refusal is transient and must not be dead-lettered. **Bound** the attempts for anything you cannot classify confidently, so a genuinely poisoned record cannot loop forever. **Divert** the deterministic failures to a dead-letter destination carrying the payload, the key, the source position and the error, then commit past the record and continue. Alert on it, because you have just chosen availability over completeness for that record. The subtlety worth raising: diverting a record breaks per-key ordering — a later update for the same key now applies while the earlier one sits in the dead-letter store, so replaying it blindly can regress the target's state.

code

text · 6 lines
text
2026-08-20T02:14:03Z ERROR consumer[orders-sink] partition=3 offset=884213 
  JsonParseException: Unexpected character ('N' (code 78)) at [Source: line 1, column 42]
2026-08-20T02:14:05Z INFO  task restarting, resuming from committed offset 884213
2026-08-20T02:14:06Z ERROR consumer[orders-sink] partition=3 offset=884213 
  JsonParseException: Unexpected character ('N' (code 78)) at [Source: line 1, column 42]
2026-08-20T02:14:08Z INFO  task restarting, resuming from committed offset 884213

go deeper

for a junior

Recognise the pattern: one bad message that fails the same way every time blocks everything behind it. Knowing that it must be set aside rather than retried forever is the level-appropriate answer.

for a middle

Explain the mechanics — why progress cannot be committed past the record, what the divert entry must carry, and how a deterministic failure differs from a transient one in its diagnostic signature.

for a senior

Show you have lived it: classify errors rather than counting retries, keep an outage from dead-lettering the entire stream, alert on divert rate and new error classes, and name the per-key ordering consequence.

for a principal

Own the policy across teams: which error classes are divertible by default, what the divert budget is before a pipeline stops instead, who is paged, and how replay is reconciled so an old message can never overwrite newer state.

## What a poison record actually is A poison record is a message whose processing fails **deterministically**. It is not slow, not unlucky, not affected by load. Every attempt produces the same exception at the same position: a payload that is not valid JSON, a field the schema says is an integer holding `"N/A"`, a value that violates a not-null constraint in the sink, an event whose enum member the code has never heard of. Because the consumer cannot process it and will not skip it, it also cannot commit progress past it. The partition or the queue stalls, lag climbs, and — if the runtime restarts the failed task — the system enters a tight crash-restart loop that also floods the logs. The diagnostic signature is the giveaway and interviewers expect you to name it: **the same offset or message id appears in the failure every single time**, and the error is identical. Compare that with a transient failure, where the failing position moves around, the error mentions the network or a downstream service, and a retry sometimes succeeds. ## Step one: classify, do not just count The most common bad implementation is "retry N times, then dead-letter, regardless of the error". It works beautifully on poison records and catastrophically during an outage. If the sink database is down, *every* record fails, every record exhausts its attempts, and the dead-letter store fills with the entire stream — correctly ordered data turned into an unordered backlog that someone must now replay by hand, often after the source's retention has expired. The dead-letter path must be reserved for errors that are **not going to succeed later**: - Deserialization and parse errors — the bytes are what they are. - Schema or contract violations — a required field is absent or the type is wrong. - Domain rejections — an amount outside its permitted range, an unknown reference key that will never resolve. And the halt-or-wait path is for errors whose success depends on something outside the record: connection failures, timeouts, authorization errors, capacity rejections. For those the right behaviour is to stop consuming and raise, not to divert — the pipeline should go down loudly rather than quietly convert an outage into permanent data displacement. A bounded attempt count still belongs in the design, but as a **safety net for misclassification**, not as the primary mechanism. ## Step two: divert with enough context to replay When a record is diverted, the dead-letter entry must carry the original payload untouched, the message key, the source position (topic, partition, offset, or queue message id), the timestamp, the consumer's version, and the exception class and message. Without the position you cannot prove where it came from; without the key you cannot reason about ordering; without the payload you cannot replay. Then commit past the record so the consumer resumes. ## Step three: make the divert loud Diverting is a decision to accept an incomplete target in exchange for a moving pipeline, and it must be visible. Emit a metric of dead-lettered records per minute, broken down by error class, and alert on two shapes: a **spike** (a schema change just made every record poisonous) and any **first appearance of a new error class**. A dead-letter destination that grows without anyone noticing is the streaming equivalent of silently dropping rows. ## The ordering problem nobody mentions until asked This is the part that separates a senior answer. In a keyed stream, records for the same key are ordered, and downstream state depends on applying them in order. Divert the record that sets `status = shipped` and the next record for that order — `status = delivered` — is applied on top of a target that never saw the shipped transition. That may be harmless if the sink is a last-write-wins upsert of full state, and actively wrong if the sink applies deltas, appends to a history, or drives a state machine. So the divert decision is per-sink-semantics, and the mitigations are concrete: for delta-style or state-machine sinks, record the diverted key so the consumer can skip or quarantine subsequent records for that key until the original is resolved; or accept the gap and, on replay, reconcile the key from the source rather than applying the stale message. A blind replay of an old dead-lettered record into a target that already holds newer state is how a replay makes things worse than the original failure. ## What good looks like Errors classified explicitly; deterministic failures diverted with full context; transient failures halting rather than diverting; attempts bounded as a backstop; a metric and alert on the divert rate and on new error classes; and a documented statement of what happens to per-key ordering when a divert occurs, with a matching replay procedure.

  • Why is a blanket attempt limit followed by diverting a dangerous design?
    Because it cannot tell a poison record from an outage. If the sink is unavailable, every record exhausts its attempts and the dead-letter store fills with the entire stream — ordered data turned into an unordered backlog, often discovered after the source's retention has expired. Classify the error first; the attempt limit is a backstop for misclassification, not the primary rule.
  • What happens to per-key ordering when you divert one record?
    It breaks for that key. A later update applies to a target that never saw the diverted transition. That is tolerable when the sink upserts full state and last-write-wins, and wrong when it applies deltas, appends history, or drives a state machine. For those sinks, track the diverted key and hold or quarantine subsequent records for it until resolved.
  • How do you replay a dead-lettered record safely once the consumer is fixed?
    Never apply it blindly on top of newer state. Either reconcile the key from the source instead of replaying the stale message, or make the write conditional on the record's own version or event timestamp so an older payload cannot overwrite a newer one. Then mark the dead-letter entry resolved so it cannot replay twice.

saying these in an interview costs you the question

  • Diverting after a fixed retry count regardless of the error class
  • Dead-lettering during a sink outage and losing the whole stream's order
  • Skipping the failing record without recording it anywhere
  • Assuming a replayed dead-letter record can always be applied as-is
  • Treating an unbounded dead-letter store as a solved problem

context