A platform team sets one default for every stream consumer: retry a failing record three times then route it to a dead-letter topic and carry on. The ledger team's consumer derives account balances from ordered per-account events. What trade-off does that default make on their behalf and what should they run instead?
Show the full answer Hide the answer
What the default buys and what it pays
Skip-and-continue buys availability of the consumer group: one malformed record cannot stop a partition, so a bug in a rarely used field does not become an outage. For a clickstream, a telemetry feed or a search-index updater that is the right default, and one lost record in ten million is invisible and harmless.
It pays with silent incompleteness, and the bill lands wherever a record's meaning depends on the records around it. A ledger is exactly that case. Balances are a fold over an ordered sequence per account: skipping event 5 and applying 6 and 7 produces a balance that is wrong by the amount of event 5 and stays wrong forever, with no error anywhere. The consumer is healthy, lag is zero, the dashboard is green, and one account is out by the value of a transaction.
Why replay does not rescue it
Dead-letter replay is an out-of-order re-application. Event 5 arrives after 7 has already been folded in, so unless every downstream effect is commutative, the state after replay differs from the state that would have existed had nothing failed. For a sum it happens to work; for "decline if balance would go negative" or "charge an overdraft fee", it does not.
The decision rule: if the effect of a record depends on the records before it, the correct response to an unprocessable record is to stop. Order-independent records can be diverted; order-dependent records cannot.
Why the other options fail
- Keep the default and alert on a non-empty dead-letter topic. This is the common compromise and it fails on timing rather than on intent. By the time the alert is read, the partition has moved on and the fold has already absorbed the gap. The alert tells you that you have a reconciliation problem, not that you avoided one.
- Raise the retry count to twenty. Retries fix transient failures. A poison record is deterministic — a schema mismatch or an unhandled enum — so twenty attempts take longer to reach the same outcome, and the extra minutes of stalled partition are pure loss.
- Dead-letter and auto-replay after ten minutes. This is the worst of the set: it looks like recovery while guaranteeing out-of-order application, and it hides the failure from whoever needed to know that a balance is provisional. It would be right for an independent record whose failure cause is expected to clear, such as a downstream service being briefly unavailable — which is a retry problem wearing a dead-letter costume.
When halting is the wrong answer
Halting is an availability decision with a real cost: the partition stops, lag grows at the arrival rate, and every downstream consumer of that account's stream stalls with it. Do not make it the platform default. Make it an option the data's owner chooses, with the blast radius already reduced — partition by account so only that account's stream halts, and page the team that owns the meaning of the data rather than the team that owns the broker.
Common weak answers
- "Exactly-once processing avoids this." Delivery semantics say nothing about a record the code cannot interpret.
- "Log it and move on; we will reconcile nightly." Reconciliation detects the divergence; it does not tell you what the missing event said, unless the record itself is still recoverable.