A stream is an ordered sequence within declared partitions, not a globally ordered table of facts.
Stream event contract: identity, version and partition key
Define one event
A receipt-submitted event needs event ID, receipt ID, occurrence time, ingestion time, schema version and producer identity. A retry must retain the same event ID if it represents the same fact. Distinguish an update from a new submission with an event type and version. Dataset contracts begin at producers, before a consumer tries to infer missing fields.
Choose the partition key
Partitioning by receipt ID keeps events for one receipt in order on a single partition, while events for different receipts can interleave. A hot key can overload a partition, so measure key distribution. Do not assume total ordering across partitions. A consumer that needs per-customer state must choose a compatible key or deliberately repartition.
Plan schema change
Additive fields with defaults can be easier to roll out than changing the meaning of an existing field. Consumers should reject unknown required versions into a quarantine path rather than guessing. Store raw event bytes or a reproducible decoder version for replay. A published contract should say which fields may be absent and how deletion events are represented.
Exercise a retry
Send the same receipt-submitted event twice with identical event ID; the downstream count should increase once. Send a correction with a new event ID and explicit superseded version; it should update the state under a declared rule. Add an event with no receipt ID and confirm it is rejected before partitioning.
Implementation
def validate_receipt_event(event):
required = {"event_id", "receipt_id", "occurred_at", "ingested_at", "schema_version"}
if required - event.keys() or not event["event_id"] or not event["receipt_id"]:
raise ValueError("missing stream identity")
if event["schema_version"] not in {2, 3}:
raise ValueError("unsupported event schema")
return event["receipt_id"]Performance and operating cost
Validation is O(1) per event and O(N) for N events. Partition-key skew can dominate end-to-end latency even when total throughput looks healthy; publish per-partition lag and key-distribution diagnostics.
Common Mistakes
- Do not assign a new event ID to a retry of the same fact.
- Do not assume one global order across partitions.
- Do not silently reinterpret a field under the same schema version.
Read next
- Event-time windows: place each event by occurrence, not arrival
- Idempotent streaming sinks: make retries safe across failures
- Backfills: rebuild history without exposing a half-written result
- Time-series calendar: distinguish missing periods from measured zeros
Continue the workflow: Product event contracts: identity, deduplication and eligibility.
