Build a versioned order table that accepts two payload schemas, prunes useful partitions, commits complete snapshots, and proves a deletion survives replay.
Project: release a versioned order lake
Fix the release contract
The source produces order events with order ID, amount in integer cents, event day, and a schema version. Version three adds currency, while version two defaults to INR. Define one current row per order ID, ordered by source position; a delete event removes that row. Keep a raw immutable batch and a rejected-record ledger so the accepted plus rejected count matches the source count.
Build two candidate partitions
Prepare day 47 and day 48 output under new file names. A late correction to an order from day 47 must rewrite or merge into that event-day partition even if it arrives on day 52. Record file sizes, min/max days, distinct keys and a checksum before commit. Partition design should make the day-bounded audit scan cheaper without creating one file per order.
Simulate a failed commit
Start from snapshot 41. Have writer A prepare a corrected order and writer B prepare a deletion against the same parent. Commit B first. Writer A must receive a stale-parent error, reload snapshot 42, apply the tombstone, and rebuild its candidate. Publishing A's old candidate would resurrect the deleted order. Atomic publication supplies the conflict boundary.
Run release checks
Decode retained version-two and version-three fixtures, check the amount unit, verify accepted plus quarantined equals extracted, and query the table through a pinned snapshot. Execute a deletion request and replay the old raw batch; the subject must stay absent. Inspect one audit query's candidate row groups and compare bytes scanned to a full scan. Record the exact snapshot ID and schema versions in the release packet.
State the outcome and rollback
Ship only when both schema fixtures decode, checksums match, the conflict test rejects stale publication, and deletion verification passes. A rollback points readers to the previous allowed snapshot only if doing so cannot expose a subject that must remain erased; otherwise repair forward. Record what evidence is still pending when backup expiry extends beyond the release window.
Implementation
source_events = [
{"order_id": "ord-47", "position": 81, "kind": "upsert", "amount_cents": 7350},
{"order_id": "ord-83", "position": 82, "kind": "upsert", "amount_cents": 4200},
{"order_id": "ord-47", "position": 84, "kind": "delete"},
]
def materialize_current(records):
current = {}
for event in sorted(records, key=lambda item: item["position"]):
if event["kind"] == "delete":
current.pop(event["order_id"], None)
else:
current[event["order_id"]] = event["amount_cents"]
return current
assert materialize_current(source_events) == {"ord-83": 4200}Performance and operating cost
Sorting N events costs O(N log N) time and O(N) space; an already ordered change log can be applied in O(N) time. File rewrite and validation cost depend on affected partitions, not only changed rows. The release packet must also price retained snapshots and the deletion ledger.
Common Mistakes
- Do not merge rows by arrival time when source position defines order.
- Do not publish a stale candidate after a concurrent deletion.
- Do not roll back to a snapshot that would expose an erased subject.
Read next
- Schema compatibility and consumer rollout
- Parquet row groups, projection and scan cost
- Partition planning, pruning and skew
- Table snapshots and atomic publication
- Compaction, retention and the replay horizon
- Deletion propagation and erasure proof
- Project: build a replayable receipt pipeline with an atomic publish gate
Continue the workflow: Project: release a distributed invoice rollup.
Continue the workflow: Project: migrate an order event contract safely.
Continue the workflow: Project: reduce delete overhead without breaking old snapshots.
