Design continuous source-vs-derived drift detection, repair queues, version markers, checksums, and semantic translation boundaries for systems that must coexist for years.
Reconciliation Jobs and Anti-Corruption Layers for Long-Lived Data Synchronization
Long-lived synchronization needs continuous proof. AtlasMart finishes the chapter by detecting drift, repairing from authority, and isolating legacy semantics.
Design continuous source-vs-derived reconciliation using counts, checksums, source versions, sampling, and repair queues.
Differentiate detection evidence from repair authority so reconciliation does not create a second uncontrolled writer.
Use an anti-corruption layer to translate legacy schemas/semantics into a canonical model and reject unknown meanings explicitly.
Treat long-lived dual systems as permanently drift-prone and prove convergence through repeated game-day/rebuild tests.
1. Streams reduce drift; they do not prove absence of drift
AtlasMart has run a new order read model beside a legacy order system for two years. Even with CDC and replay, defects can accumulate: a connector filter was wrong for one week, one target rejected a malformed value, a schema migration changed defaults, or an operator manually edited a derived record. A reconciliation job compares authoritative source evidence with a derived copy and produces explicit drift findings.
Reconciliation is not merely “compare row counts.” Counts can match while values differ. Useful evidence layers include counts by partition/time range, checksums of canonicalized records, source-version/watermark comparisons, targeted samples, and full key/value comparison for high-risk subsets.
2. Detection and repair are different privileges
A reconciliation process may have broad read access to compare
systems, but its write authority should be carefully bounded. A
repair queue turns findings into explicit
idempotent actions such as
UPSERT o-2 from source version 6 or
DELETE orphan o-9. The repair worker verifies that
the source has not changed again before writing.
A job powerful enough to rewrite every derived store can become a mass-corruption mechanism. Separate read-only detection, reviewed/automated repair policy, credentials, audit logging, rate limits, and rollback.
3. Anti-corruption layers protect semantics, not only syntax
An anti-corruption layer is a translation
boundary between two models. The legacy order service may store
state='P' and cents as integers; the new canonical
API expects status='PAID' and a money
representation. The layer maps names, types, units, missing
values, enum meaning, and identity. Critically, unknown legacy
state 'Z' should fail/quarantine rather than be
guessed as PAID.
This boundary prevents legacy quirks from spreading through every new consumer. It also creates a place to version transformations and test semantic equivalence.
4. Deliberately wrong approach: run one migration verification and declare the systems synchronized forever
Long-lived systems continue changing. New code, backfills, operator corrections, retention, retries, partial outages, and schema evolution can reintroduce drift tomorrow. A one-time migration script proves one historical point, not an ongoing invariant.
The safer design schedules reconciliation according to business risk, exposes drift rate/age, and keeps repair/rebuild procedures executable. High-value records such as payment state or entitlement may require continuous or near-real-time cross-checks; low-value analytics projections can use sampled/batched reconciliation.
5. AtlasMart lab: detect stale/missing/orphan rows and translate legacy semantics
Python 3.13+ standard library only. SHA-256 is used only as a deterministic teaching checksum. No production database scan, destructive repair, credentials, or network access occurs.
import hashlib, json
source = {
"o-1":{"status":"PAID","version":4,"total":120},
"o-2":{"status":"SHIPPED","version":6,"total":80},
"o-3":{"status":"CANCELLED","version":3,"total":45},
}
derived = {
"o-1":{"status":"PAID","source_version":4,"total":120},
"o-2":{"status":"PAID","source_version":5,"total":80}, # stale
"o-9":{"status":"PAID","source_version":1,"total":999}, # orphan
}
def digest(d):
payload=json.dumps(d,sort_keys=True,separators=(',',':')).encode()
return hashlib.sha256(payload).hexdigest()[:12]
print("COUNTS / CHECKSUMS")
print("source count",len(source),"derived count",len(derived))
print("source digest",digest(source),"derived digest",digest(derived))
repair=[]
for oid,truth in source.items():
copy=derived.get(oid)
if copy is None or copy.get('source_version') != truth['version']:
repair.append(("UPSERT",oid,truth))
for oid in set(derived)-set(source):
repair.append(("DELETE",oid,None))
print("repair queue:", [(op,oid) for op,oid,_ in repair])
for op,oid,truth in repair:
if op=="UPSERT":
derived[oid]={"status":truth['status'],"source_version":truth['version'],"total":truth['total']}
else:
derived.pop(oid,None)
print("post-repair equal IDs:", sorted(source)==sorted(derived))
print("versions:", {k:derived[k]['source_version'] for k in sorted(derived)})
print("\nANTI-CORRUPTION LAYER")
legacy_rows=[
{"id":"x1","state":"P","amount_cents":2500},
{"id":"x2","state":"X","amount_cents":900},
]
state_map={"P":"PAID","X":"CANCELLED"}
def translate(row):
if row['state'] not in state_map:
raise ValueError("unknown legacy state")
return {"order_id":row['id'],"status":state_map[row['state']],"total":row['amount_cents']/100}
print("translated:", [translate(r) for r in legacy_rows])
print("\nCONTINUOUS, NOT ONE-TIME")
print("a migration that was correct yesterday can drift tomorrow; schedule comparison + repair with bounded write authority")
Counts and checksums differ. The repair queue identifies stale
o-2, missing o-3, and orphan
o-9; after repair, derived IDs and source
versions match. The anti-corruption layer maps two documented
legacy states into the canonical representation.
6. Production judgment and chapter synthesis
Partition large reconciliation jobs by stable key/range/time bucket, rate-limit reads, and avoid turning validation into a production outage. Canonicalize data before checksumming so irrelevant serialization order does not generate false drift. Record comparison run ID, source snapshot/watermark, target watermark, counts, mismatch samples, repair actions, completion state, and operator approval when required.
Chapter 20's full mechanism is now explicit: capture committed changes or publish semantic events; make publication intent durable; transport with partition-scoped ordering and replay; build independent derived views with freshness evidence; then continuously reconcile because every asynchronous copy can drift. At-least-once delivery is a recoverability strategy, not a promise that duplicates never exist.
Chapter 21 turns to the security boundary around all of this data movement: identities, least privilege, encryption, tenant isolation, audit, retention, and administrative interfaces.
Check your understanding
- Why are equal row counts insufficient reconciliation evidence?
- What is a repair queue for?
- What does an anti-corruption layer translate?
- Why should unknown legacy values be quarantined rather than guessed?
- Why must reconciliation be continuous for long-lived dual systems?
Review the answers
1. The same keys can contain different values/versions, and missing plus orphan records can cancel out numerically.
2. It turns detected drift into explicit, auditable, idempotent correction actions rather than ad-hoc writes.
3. Schema and semantic differences such as names, units, types, identities, enum meanings, defaults, and error handling.
4. Guessing invents business meaning and can silently corrupt canonical state.
5. New failures, schema changes, operations, retries, and bugs can recreate drift after an initially correct migration.
References
Foundational claims use primary specifications/research or current official documentation where practical. Product references are optional implementation anchors; the mandatory labs are vendor-neutral.
- PostgreSQL 18 logical decoding concepts — Source position, replay, and duplicate behavior relevant to reconciliation boundaries.
- Debezium 3.6.1.Final release notes — Current stable CDC connector snapshot including data-integrity and credential-leak fixes.
- Transactional Outbox pattern — Durable publication-intent pattern used earlier in the chapter.
- Microsoft anti-corruption layer pattern — Architecture pattern for semantic isolation between legacy and modern models.