Build an hourly receipt-submission stream with explicit event-time windows, duplicate handling, late corrections and batch reconciliation.
Project: build a replay-safe receipt event dashboard
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
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 countsPerformance 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.
