A consumer can see an event again after a crash or replay; its externally visible effect needs a stable identity and commit rule.
Idempotent streaming sinks: make retries safe across failures
Locate the failure gap
If a consumer writes an hourly count and crashes before committing its input position, it will process the same event again. At-least-once delivery is often workable when the sink deduplicates by event ID or writes a versioned aggregate idempotently. A producer-side duplicate guarantee does not automatically cover a nontransactional downstream database. Event identity is the key.
Choose the sink operation
A blind increment is unsafe on replay. An upsert keyed by window, metric and version can replace the same result, while an event ledger with a unique event ID can guard incremental updates in one transaction. Document how corrections supersede previous values and how a deleted receipt is reversed. Avoid claiming end-to-end exactly-once behavior without checking every boundary.
Keep replay deterministic
A replay needs the same schema decoder, event selection rule and aggregation version as the original run, or a deliberate migration. Store input offsets and code version with checkpoints. Backfill publication can rebuild into a separate snapshot and switch readers once counts reconcile.
Test a crash
Process a three-event fixture, crash after writing the second result but before recording progress, then replay from the previous checkpoint. The final count must match a clean run. Add a correction and verify it changes the intended window once, including after a second replay. Compare both event ledger and visible aggregate.
Implementation
def apply_once(event, seen_event_ids, counts_by_window):
event_id = event["event_id"]
if event_id in seen_event_ids:
return False
window_key = event["window_key"]
seen_event_ids.add(event_id)
counts_by_window[window_key] = counts_by_window.get(window_key, 0) + 1
return TruePerformance and operating cost
A hash-set check is O(1) expected per event but an unbounded set uses O(N) memory. Production deduplication needs a durable bounded identity store or a sink transaction; the in-memory fixture demonstrates semantics, not crash durability.
Common Mistakes
- Do not use blind increments with replayable input.
- Do not equate producer idempotence with an end-to-end sink guarantee.
- Do not discard replay metadata before the retention window ends.
