Move a payment mart from one amount field to a currency-aware contract while preserving old readers, late events and a reversible release pointer.
Project: migrate a live payment amount contract
Prepare the contract
Create 470 payment records under the old schema, with separate settled and refunded events. Record schema IDs, readers, source positions and the maximum supported replay interval. Add amount_minor_units and currency without removing amount_cents. Require dual-write equality for existing currency records and quarantine records lacking a trustworthy currency. The staged migration prevents one deployment from forcing all consumers to change at once.
Repair history in bounded units
Split historical data into event-date partitions and pin a source snapshot for each repair. Persist partition progress and candidate output generation. Inject a live correction during the repair; its later source position must win over the old snapshot. Rerun one partition to prove idempotence. The candidate remains hidden until all partitions needed by the reader contract have passed.
Reconcile old and new paths
Calculate totals, refunded amounts, primary-key sets, duplicate counts and null rates by currency and date from the same source position. Introduce a duplicated refund so a global total comparison is insufficient. The cutover gate must fail and report the first differing key. Repair it, rerun the check and retain both failed and successful evidence.
Cut over and roll back
Build the new table and serving index as one candidate release. Move a canary reader, then all approved readers, using versioned pointers. Simulate an index build failure before publication and confirm the old reader still answers. After a successful cutover, revert the pointer once and verify the old representation remains valid during the overlap window.
Submit release evidence
Deliver the reader inventory, schema changes, dual-write comparison, partition checkpoints, pinned source positions, mismatch report, fixed reconciliation, old and new pointer IDs, rollback result and delayed old-schema replay. Do not report completion until the replay horizon closes and owner review confirms that retiring amount_cents will not break a dormant export.
Implementation
payments = [("pay-47", 2375), ("pay-48", -125), ("pay-49", 6400)]
def materialize_amounts(records):
return {payment_id: {"amount_cents": cents,
"amount_minor_units": cents, "currency": "USD"}
for payment_id, cents in records}
candidate = materialize_amounts(payments)
assert sum(row["amount_minor_units"] for row in candidate.values()) == 8650
assert all(row["amount_cents"] == row["amount_minor_units"]
for row in candidate.values())Performance and operating cost
The sample materialization uses O(N) time and O(N) state. In a real lake or warehouse, the migration pays for dual writes, candidate snapshots, index rebuilds and partition-level comparison scans. A bounded repair budget limits blast radius and makes retries affordable; premature removal of the old field saves bytes while risking delayed consumers and historical replay.
Common Mistakes
- Do not drop the old field before delayed readers and exports have moved.
- Do not let a duplicated refund pass because another error cancels its total.
- Do not publish a table whose serving index build failed.
