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

Project: build a replay-safe receipt event dashboard

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

Build an hourly receipt-submission stream with explicit event-time windows, duplicate handling, late corrections and batch reconciliation.

Prepare the event fixture

Create submission events with stable event IDs, receipt IDs, occurrence and ingestion times, schema versions and partition keys. Include a retry of the same event, an out-of-order submission, an unsupported schema and a correction after the first hourly output. The event contract must state which records count as submissions.

Compute and publish

Assign each event to a half-open one-hour event-time window. Deduplicate by event ID, record a provisional count and route events after the watermark under a declared grace policy. Publish window ID, version, count, watermark and last update time. A replayed event must not increment the visible count twice. Sink behavior must survive the fault fixture.

Reconcile independently

Calculate counts from a fixed raw-event snapshot with a separate batch pass and compare every finalized window. Explain discrepancies by event ID and correction version. Simulate a consumer crash after writing but before checkpointing; the final output must match the clean run. Measure state size, backlog and correction volume during the run.

Submit a handoff

Provide source events, schema definition, dedup ledger or sink transaction plan, window results, replay trace, parity report and dashboard state example. Mark an idle source stale rather than displaying a confirmed zero. Include a recovery procedure that rebuilds into a separate snapshot and switches only after reconciliation.

Implementation

python
def receipt_hourly_counts(events):
    counts = {}
    seen_ids = set()
    for event in events:
        event_id = event["event_id"]
        if event_id in seen_ids:
            continue
        seen_ids.add(event_id)
        if event["type"] != "submitted":
            continue
        hour_start = (event["occurred_minute"] // 60) * 60
        counts[hour_start] = counts.get(hour_start, 0) + 1
    return counts

Performance and operating cost

A pass over N events is O(N) expected time and O(U + W) memory for U unique IDs and W windows. A production system must bound durable dedup and window state, and its batch parity check must use the same event-selection and correction rules.

Common Mistakes

  • Do not count a retried event twice.
  • Do not publish a provisional window as final without a watermark rule.
  • Do not claim replay safety from an in-memory set alone.

Read next

ai-data
streaming-analytics-project
Storage details