Source-to-target reconciliation at lakehouse scale
Overview
A distributed data platform is a chain of hand-offs. Operational databases emit changes, a CDC connector publishes them to Kafka, PySpark jobs clean and reshape them, a lakehouse table format stores the result, and Trino serves it to analysts, dashboards and increasingly to AI systems. Each hand-off can lose, duplicate or quietly alter data, and none of them will raise an error when it does.
This is a reference architecture for answering one question with evidence: did the data that left the source arrive in the target, complete and correct? It covers reconciliation (counts, keyed diffs, checksums, partition-level comparison), schema and transformation validation, the classic integrity checks, and the monitoring layer that keeps watching after release.
In 2023 I wrote a three-part series on big data testing: Part 1 on why it matters, Part 2 on practices and frameworks such as Great Expectations and PyDeequ, and Part 3 on wiring Great Expectations to Trino. Those pieces answer “is this dataset well-formed?”. This one goes a level deeper and asks “is this dataset the same data as the source, after the transformations we intended and no others?”. That second question is where most production incidents actually live.
Problem statement
Single-table expectations catch malformed data. They do not catch missing data that looks fine. A partition that silently loaded 92% of its rows passes every null, range and uniqueness check. So does a join that dropped orders with an unmatched customer key, or a CDC consumer that skipped a batch of deletes after a rebalance.
The typical failure modes in this kind of pipeline are:
- Loss. Records dropped by a failed micro-batch, a filter with an unintended predicate, an inner join where a left join was meant, or Kafka retention expiring before a lagging consumer caught up.
- Duplication. At-least-once delivery replayed after a restart, or an append where an upsert was intended.
- Mutation. Type coercion (decimal to double), timezone shifts, truncated strings, or a transformation rule implemented differently from its specification.
- Drift. A source column added, renamed or retyped, so the pipeline either fails loudly (good) or maps nulls into a column for weeks (bad).
- Staleness. Everything is correct, but six hours late.
Manual ETL validation, where an engineer writes a few comparison queries before sign-off, does not scale across hundreds of tables and continuous loads. In my practice I built an LLM-powered data validation framework that automated more than 70% of manual ETL checks (the practice case study has the background). The design below is the reconciliation and validation layer that such automation needs underneath it: deterministic, repeatable checks with clear severities, so that automation generates and runs rules rather than inventing judgements.
Engineering objectives
- Completeness. Every in-scope source record is present in the target exactly once, per key and per partition.
- Accuracy. Values in the target match the source after applying the documented transformation, and only that transformation.
- Schema integrity. Structural changes at the source are detected before they reach consumers.
- Consistency across systems. Source database, Kafka topic, lakehouse table and serving views agree at a defined point in time.
- Scalability. Validation cost grows with change volume, not total table size, and runs where the data lives.
- Actionability. Every failed check has a severity, an owner and a defined consequence: block, quarantine or alert.
Solution architecture
- 01Source DB
- 02Debezium CDC
- 03Kafka topics
- 04PySpark transform
- 05Iceberg tables
- 06Trino serving
The validation layer sits beside this path rather than inside it. It reads from every stage, stores results in its own tables, and talks back to the orchestrator only through gates.
01 Capture
- Source snapshots and counts
- Kafka offsets and lag
- Schema registry versions
02 Validate
- Rule catalogue (YAML)
- Trino reconciliation SQL
- PySpark checks for heavy diffs
- Great Expectations suites
03 Decide
- Severity policy
- Airflow gates
- Quarantine tables
04 Observe
- Results store
- Anomaly baselines
- Alerts and dashboards
Two design choices shape everything else.
Trino as the reconciliation engine. Because Trino federates catalogs, a single query can read the PostgreSQL source and the Iceberg target and compare them. That removes a whole class of “export both sides and diff in Python” scripts, and it pushes the work to an engine built for scans and aggregations.
PySpark for the heavy diffs. When a full keyed comparison over billions of rows is needed, for example after a backfill or a transformation rewrite, it runs as a Spark job close to the storage, not through an interactive query cluster.
Technical approach
Schema and metadata validation
Schema drift is cheapest to catch first. I compare information_schema.columns across catalogs and fail on any difference that the mapping does not explain.
WITH src AS ( SELECT column_name, data_type FROM pg_sales.information_schema.columns WHERE table_schema = 'public' AND table_name = 'orders'),tgt AS ( SELECT column_name, data_type FROM lake.information_schema.columns WHERE table_schema = 'sales' AND table_name = 'orders')SELECT coalesce(s.column_name, t.column_name) AS column_name, s.data_type AS source_type, t.data_type AS target_type, CASE WHEN t.column_name IS NULL THEN 'missing_in_target' WHEN s.column_name IS NULL THEN 'extra_in_target' ELSE 'type_changed' END AS driftFROM src sFULL OUTER JOIN tgt t ON s.column_name = t.column_nameWHERE s.column_name IS NULL OR t.column_name IS NULL OR s.data_type <> t.data_type;Type names differ across connectors (character varying against varchar), so in practice the comparison joins through a small mapping table of accepted equivalences. Extra target columns that the pipeline adds, such as _ingested_at, go on an allow-list. Anything else is a finding.
Row counts, per partition
A global row count is a smoke test. It hides offsetting errors: 1,000 rows lost in one day and 1,000 duplicated in another sum to zero. Reconcile per partition, usually per business date.
WITH src AS ( SELECT CAST(created_at AS date) AS d, count(*) AS n FROM pg_sales.public.orders WHERE created_at >= DATE '2026-09-01' AND created_at < DATE '2026-10-01' GROUP BY 1),tgt AS ( SELECT order_date AS d, count(*) AS n FROM lake.sales.orders WHERE order_date >= DATE '2026-09-01' AND order_date < DATE '2026-10-01' GROUP BY 1)SELECT coalesce(s.d, t.d) AS d, s.n AS source_rows, t.n AS target_rows, coalesce(t.n, 0) - coalesce(s.n, 0) AS deltaFROM src s FULL OUTER JOIN tgt t ON s.d = t.dWHERE coalesce(s.n, -1) <> coalesce(t.n, -1)ORDER BY d;For Iceberg targets, the $partitions metadata table already holds record_count per partition, so the target side of a frequent count check can be answered from metadata without scanning data files. I use that for the every-load check and keep the full scan for the nightly run.
Aggregate checksums
Counts prove quantity, not content. An order-insensitive checksum over a canonical string of each row proves content cheaply. Trino’s checksum aggregate is order-insensitive, which matters because source and target will never return rows in the same order.
WITH src AS ( SELECT CAST(created_at AS date) AS d, checksum(concat_ws('|', CAST(order_id AS varchar), coalesce(CAST(customer_id AS varchar), '~'), coalesce(CAST(CAST(amount AS decimal(18,2)) AS varchar), '~'), coalesce(status, '~'))) AS h, sum(CAST(amount AS decimal(38,2))) AS amount_total FROM pg_sales.public.orders GROUP BY 1),tgt AS ( SELECT order_date AS d, checksum(concat_ws('|', CAST(order_id AS varchar), coalesce(CAST(customer_id AS varchar), '~'), coalesce(CAST(amount AS varchar), '~'), coalesce(status, '~'))) AS h, sum(amount) AS amount_total FROM lake.sales.orders GROUP BY 1)SELECT s.d, s.amount_total, t.amount_totalFROM src s JOIN tgt t ON s.d = t.dWHERE s.h <> t.h OR s.amount_total <> t.amount_total;Two details carry most of the value. First, canonicalise before hashing: cast decimals to the target precision, normalise timestamps to UTC, and replace nulls with a sentinel so that a null and an empty string do not collide. Second, apply the same transformation on the source side that the pipeline applies. If the pipeline upper-cases status, the source expression must too, otherwise the checksum compares two different things and fails forever.
A mismatched partition checksum does not say which row is wrong. It tells you where to spend a keyed diff.
Keyed diffs
Inside a failing partition, a full outer join on the business key classifies every discrepancy.
SELECT coalesce(s.order_id, t.order_id) AS order_id, CASE WHEN t.order_id IS NULL THEN 'missing_in_target' WHEN s.order_id IS NULL THEN 'unexpected_in_target' ELSE 'value_mismatch' END AS finding, s.amount AS src_amount, t.amount AS tgt_amount, s.status AS src_status, t.status AS tgt_statusFROM (SELECT * FROM pg_sales.public.orders WHERE CAST(created_at AS date) = DATE '2026-09-14') sFULL OUTER JOIN (SELECT * FROM lake.sales.orders WHERE order_date = DATE '2026-09-14') t ON s.order_id = t.order_idWHERE s.order_id IS NULL OR t.order_id IS NULL OR CAST(s.amount AS decimal(18,2)) IS DISTINCT FROM t.amount OR s.status IS DISTINCT FROM t.statusLIMIT 1000;IS DISTINCT FROM is deliberate. A plain <> returns null when either side is null, and the row silently drops out of the result.
For whole-table diffs at scale, the same logic runs in PySpark, hashing rows first so the shuffle carries a key and a hash rather than every column.
from pyspark.sql import functions as F
COLS = ["customer_id", "amount", "status"]
def canonical_hash(df): parts = [F.coalesce(F.col(c).cast("string"), F.lit("~")) for c in COLS] return df.select("order_id", F.sha2(F.concat_ws("|", *parts), 256).alias("h"))
src = canonical_hash( spark.read.jdbc(src_url, "public.orders", properties=jdbc_props, column="order_id", lowerBound=1, upperBound=max_id, numPartitions=64) .withColumn("amount", F.col("amount").cast("decimal(18,2)")))tgt = canonical_hash(spark.table("lake.sales.orders"))
diff = (src.alias("s") .join(tgt.alias("t"), "order_id", "full_outer") .select("order_id", F.when(F.col("t.h").isNull(), "missing_in_target") .when(F.col("s.h").isNull(), "unexpected_in_target") .when(F.col("s.h") != F.col("t.h"), "value_mismatch") .alias("finding")) .where(F.col("finding").isNotNull()))
summary = diff.groupBy("finding").count().collect()diff.limit(10_000).write.mode("overwrite").saveAsTable("dq.orders_diff_sample")The partitioned JDBC read keeps the source database from serving one enormous query. Even so, full diffs against an operational database belong in a read replica or a snapshot, never the primary.
Transformation validation
Reconciliation proves the target matches the source through the intended transformation. Transformation tests prove the intended transformation is the right one. I treat each transformation rule as a specification with test cases: a small, hand-built input set covering edge cases (nulls, boundary dates, currency rounding, late-arriving dimension keys) with expected output, run in CI against the PySpark job using a local Spark session. These tests are fast, deterministic and catch logic errors before any production data is involved. Reconciliation then catches the errors that only appear with real data volume and real disorder.
Integrity checks
The classic checks still apply, now expressed as rules with severities:
- Duplicates on the business key:
count(*) - count(DISTINCT order_id)per partition must be zero. - Nulls in required columns, and null rate in optional ones compared with a baseline.
- Referential integrity: orphaned foreign keys, found with an anti-join from
orderstocustomers. In a lakehouse fed by independent pipelines, a small orphan rate can be legitimate while the dimension catches up, so this check often carries a grace window rather than a zero threshold.
Consistency across distributed systems
“Source equals target” only has meaning at a defined point in time. A source that is still accepting writes will never exactly match a target that lags by a few seconds. I pin both sides:
- On the source, reconcile only partitions that are closed (yesterday and older), or read from a snapshot at a recorded log position.
- On Kafka, record the consumer group offsets the batch committed.
- On Iceberg, record the snapshot ID the load produced and query it with
FOR VERSION AS OF, so that later loads do not move the target under the check.
That triple, source position, Kafka offsets and target snapshot, is stored with every reconciliation result. When a check fails, it is reproducible.
Rule catalogue and severity
Rules live in version-controlled YAML, reviewed like code. Each has an owner and a severity, and severity determines the pipeline’s behaviour.
table: lake.sales.ordersowner: sales-data-engrules: - id: orders.partition_count_match type: reconcile_count source: pg_sales.public.orders partition: order_date tolerance: 0 severity: critical - id: orders.partition_checksum_match type: reconcile_checksum columns: [order_id, customer_id, amount, status] severity: critical - id: orders.customer_fk type: referential references: lake.sales.customers.customer_id max_orphan_rate: 0.001 grace_hours: 6 severity: high - id: orders.row_volume type: anomaly_volume baseline_days: 28 severity: medium| Severity | Example rules | Pipeline action | Who is notified |
|---|---|---|---|
| Critical | Count or checksum mismatch, duplicate keys, breaking schema drift | Block publish; keep previous snapshot live | Owning team, paged |
| High | Orphan rate above threshold, required-column nulls | Publish to quarantine; downstream jobs wait | Owning team, ticket |
| Medium | Volume or freshness anomaly, distribution shift | Publish; annotate the dataset as suspect | Channel alert |
| Low | Additive schema change, optional-column null drift | Publish; log for weekly review | Digest only |
The table is the contract between data producers and consumers. Without it, every failure becomes a negotiation at 2 a.m.
Anomaly detection and monitoring
Reconciliation answers “did we move it correctly?”. It cannot answer “is the source itself sane today?”. If the source system only wrote half its usual orders, source and target will agree perfectly. Volume, freshness and distribution monitors cover that gap. The design of those baselines is the subject of the companion article, From expectations to a production data-quality system.
Validation strategy
Sampling versus full validation
| Approach | Catches | Misses | Use for |
|---|---|---|---|
| Global counts | Gross loss | Offsetting errors, content changes | Every load, as a smoke test |
| Partition counts (metadata) | Partition-level loss and duplication | Content changes | Every load |
| Partition checksums | Any content change, per partition | Which row | Every load for critical tables; nightly otherwise |
| Keyed diff on a failing partition | Exact rows and columns | Nothing within scope | On checksum failure |
| Random sample diff | Systematic transformation errors | Rare, localised errors | Very large tables, between full runs |
| Full keyed diff | Everything | Nothing | Backfills, migrations, transformation rewrites |
Sampling is attractive and often misleading. A 1% random sample has a good chance of finding an error that affects many rows, and almost no chance of finding the 40 rows lost from one partition. Checksums give full coverage for the cost of an aggregation, which is why I prefer them as the default and treat sampling as a supplement for value-level checks that are too expensive to run everywhere.
Pushdown
Every check is written to run where the data is. Trino aggregates on both catalogs and pushes predicates and, depending on the connector, aggregations into the source database. The validation code moves summaries, not rows. A reconciliation framework that pulls both tables into a Python process will work in the demo and fall over in month two.
Incremental and CDC validation
For CDC-fed tables, full-table reconciliation each load is wasteful. Validation follows the change volume instead: reconcile the partitions touched by the batch, compare change counts by operation type between the Kafka batch and the lakehouse commit, and run the full checksum on a slower cadence. Ordering, deletes, tombstones and idempotent MERGE each need their own tests, which are covered in Validating CDC pipelines.
Testing the pipeline itself
The validation layer is tested like any other system. I keep a set of fault-injection fixtures: a batch with a dropped partition, a duplicated micro-batch, a retyped column, a shifted timezone, an orphaned key. Each fixture must trigger a specific rule at a specific severity. If a refactor of the rule engine stops a fixture from firing, CI fails. This is the step most teams skip, and it is the only evidence that a green dashboard means anything.
Metrics and measurement
I measure the validation system on four things, all computed from its own results store:
- Coverage: share of in-scope tables with at least one critical reconciliation rule, and share of critical columns covered by a checksum.
- Detection: fault-injection fixtures caught, by rule, per release of the framework.
- Time to detect: time from a load commit to the first failing result for that load.
- Noise: share of alerts closed as “expected” or “no action”, tracked per rule. A rule with persistent noise is retuned or demoted, not ignored.
Targets for these depend on the platform and are set with the data owners. I avoid publishing a single “data quality score”. It averages away exactly the failures that matter.
Challenges and trade-offs
- Load on the source. Reconciliation queries compete with production traffic. Use replicas, snapshots or off-peak windows, and partition the reads.
- Moving targets. Without pinned positions and snapshots, results flap. Pinning adds bookkeeping but makes every result reproducible.
- Legitimate divergence. Soft deletes, GDPR erasure and filtered records mean the target should not match the source exactly. Those rules belong in the transformation specification and are applied to the source side of the comparison, not handled as exceptions.
- Cost of full diffs. They are the only proof for migrations, and they are expensive. Schedule them for events (backfills, rewrites), not as a daily habit.
- Ownership. A rule without an owner becomes an alert nobody reads. The catalogue rejects rules without one.
Outcome and lessons learned
This reference architecture turns “the data looks fine” into a set of reproducible statements: counts and checksums match per partition at a recorded source position and target snapshot; schema differences are known and accepted; integrity rules pass at their declared severities; and monitors show volume, freshness and distribution within their baselines.
The main lessons from my practice are about order and economics. Start with schema and partition counts because they are cheap and catch the loudest failures. Add checksums before sampling, because they give full coverage for the price of an aggregation. Reserve keyed diffs for localising failures and for one-off events. And put severity and ownership in the catalogue from the first rule, because that is what lets automation, including the LLM-powered validation framework described in the practice case study, add rules safely without adding noise.
Further reading
- Trino documentation: aggregate functions, including checksum
- Trino Iceberg connector: metadata tables and time travel
- Debezium PostgreSQL connector documentation
- Great Expectations documentation
- From expectations to a production data-quality system
- Validating CDC pipelines
- Letting language models check the data
- Medium series: Part 1, big data testing, Part 2, practices and frameworks, Part 3, Great Expectations