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

Stream event contract: identity, version and partition key

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

A stream is an ordered sequence within declared partitions, not a globally ordered table of facts.

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

python
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

Continue the workflow: Product event contracts: identity, deduplication and eligibility.

ai-data
streaming-analytics
Storage details