Release a new order-event schema while preserving older readers, deletion semantics and atomic warehouse snapshots under concurrent writes.
Project: migrate an order event contract safely
Freeze the current contract
Collect supported event schema IDs, retained source positions, consumer versions and one authorized replay fixture per version. Document event ID, namespaced order key, event time and delete operation separately. State the last checkpoint from which each consumer may restart. Transitive compatibility must cover that replay horizon.
Introduce the new field
Add a nullable fulfillment region with a documented unknown value. Register the proposed schema only after old payloads decode in the candidate reader and old readers accept newly written events under the selected compatibility mode. If either direction fails, stage a dual-reader or new-topic migration instead of changing a registry setting to hide the failure.
Replay changes deterministically
Feed an upsert, an order correction, a delete and a delayed retransmission. Apply source positions so the old upsert cannot resurrect the deleted order. Retain the event ID and schema ID in the audit row. Envelope rules make retry and delete behavior visible to the warehouse loader.
Publish through a conflict gate
Build a candidate mart snapshot from pinned input positions. Start an unrelated append writer from the same base version, commit it first, then force the migration writer to revalidate. It must either recompute against the new snapshot or fail cleanly; it must never replace the append. Verify one current pointer and no reader-visible candidate files.
Hand over proof
Deliver the old and new schema contracts, compatibility result against all supported versions, replay output, deletion check, competing-writer transcript and committed manifest. Include a rollback plan that restores the old reader without losing the new producer events, or explicitly state why rollback requires a compensating migration.
Implementation
order_events = [
{"id": "evt-71", "order": "ord-71", "position": 71, "op": "upsert", "cents": 6300},
{"id": "evt-72", "order": "ord-71", "position": 72, "op": "delete"},
{"id": "evt-71", "order": "ord-71", "position": 71, "op": "upsert", "cents": 6300},
]
def materialize(records):
latest, values = {}, {}
for event in records:
key = event["order"]
if event["position"] <= latest.get(key, -1):
continue
latest[key] = event["position"]
if event["op"] == "delete":
values.pop(key, None)
else:
values[key] = event["cents"]
return values
assert materialize(order_events) == {}Performance and operating cost
Replay is O(N) expected time and O(K) entity state for N events and K order keys. Registry checks, candidate table files and concurrent validation add release cost. The project accepts that overhead to prevent a lower-cost migration from making historical events unreadable or dropping a competing writer’s valid rows.
Common Mistakes
- Do not change compatibility policy merely to pass one release.
- Do not let a delayed event resurrect a deleted order.
- Do not publish a candidate snapshot before conflict validation.
Read next
- Schema registry and transitive compatibility
- Event envelopes, identity and tombstones
- Optimistic table commits and write conflicts
- Schema compatibility and consumer rollout
- Table snapshots and atomic publication
Continue the workflow: Project: migrate a live payment amount contract.
