An event envelope carries routing and replay metadata around the business payload; a tombstone makes a deletion explicit to downstream consumers.
Event envelopes, identity and tombstones
Separate three identities
An order event needs a stable event ID for deduplication, an entity key for partitioning and the source position for ordered replay. These values answer different questions. Two updates to one order share the entity key but must not share the event ID. A retransmission of the same update keeps its event ID. CDC replay uses source order to settle conflicting updates.
Carry a versioned payload
The envelope should include event type, schema identifier, source namespace, event time and publication time. A consumer can then decode the payload without inferring its shape from a nullable field. Keep business identifiers and raw payload only where classification allows them. Access boundaries must cover dead-letter storage because malformed envelopes often contain the original sensitive event.
Make deletion unambiguous
A null amount is not an order deletion. Emit an explicit delete operation or a key plus tombstone according to the stream contract. The consumer removes or marks the entity at the ordered source position, then records the delete in a replay manifest. If an older update arrives later, source-order checks must prevent resurrection. Erasure proof extends this into retained snapshots and replicas.
Bound duplicate memory
Retaining every event ID forever makes state grow with lifetime throughput. Define a deduplication horizon at least as long as the supported retry and replay window, or use source position and per-entity ordering where that is sufficient. State what happens to a duplicate outside the horizon. An event log retained for months cannot rely on a one-day ID cache for exact historical replay.
Test reorder and compaction
Send update A, delete B, retransmit A and then update C for the same entity. The final state should reflect C only if C is a legitimate later source operation. When compacting a keyed log, preserve tombstone retention long enough for lagging readers; removing it too early can leave a stale copy alive downstream.
Implementation
events = [
{"event_id": "evt-47", "order_id": "ord-47", "position": 47, "op": "upsert", "cents": 4200},
{"event_id": "evt-48", "order_id": "ord-47", "position": 48, "op": "delete"},
{"event_id": "evt-47", "order_id": "ord-47", "position": 47, "op": "upsert", "cents": 4200},
]
def replay_order(records):
state, latest, seen = {}, {}, set()
for event in records:
if event["event_id"] in seen:
continue
seen.add(event["event_id"])
key = event["order_id"]
if event["position"] <= latest.get(key, -1):
continue
latest[key] = event["position"]
if event["op"] == "delete":
state.pop(key, None)
else:
state[key] = event["cents"]
return state
assert replay_order(events) == {}Performance and operating cost
A keyed replay is O(N) expected time for N events and uses O(E + K) state for retained event IDs and K entity positions. A compacted topic reduces old payload volume but cannot replace the source audit log when exact change history matters. Tombstone retention and deduplication state must cover the slowest supported consumer.
Common Mistakes
- Do not use the entity key as the unique event ID.
- Do not interpret a nullable business value as a deletion.
- Do not expire tombstones before a lagging consumer can observe them.
Read next
- Change data capture: apply updates and deletes by source order
- Deletion propagation and erasure proof
- Schema registry and transitive compatibility
- Project: migrate an order event contract safely
- Replay manifests and audit trails
Continue the workflow: Replica version erasure and restore guards.
