Skip to content
AITroveRead. Build. Understand.
Make this comfortable

Idempotent streaming sinks: make retries safe across failures

Last updated: 5 Oct 20265 min read
tutorial
IntermediateBy AITrove Editorial

A consumer can see an event again after a crash or replay; its externally visible effect needs a stable identity and commit rule.

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

python
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 True

Performance 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.

Read next

ai-data
streaming-analytics
Storage details