Skip to content
Satya Prakash Solanki

Batch pipelines are comparatively easy to validate. The input is a file or a table, the output is a table, and you can compare the two. Change-data-capture pipelines are harder because the input is a stream of events about the source, not the source itself. The target is correct only if every event was applied, in the right order, exactly once in effect, including the events that remove data.

This part of the series collects the checks I use for incremental and CDC pipelines built from a Debezium-style connector, Kafka, PySpark and a lakehouse table queried through Trino. It is marked as working: the principles are stable, but the details vary with connector settings and table formats, and I expect to refine them.

What can go wrong

A CDC pipeline has more failure modes than a batch one, and most of them produce a target that looks healthy.

  • Reordering. Kafka guarantees order within a partition, not across a topic. If events for one key land in different partitions, or a job sorts by the wrong field, an older update can overwrite a newer one.
  • Duplicates. Connectors and consumers are typically at-least-once. A restart replays events that were already applied.
  • Late events. An event arrives after the batch that should have contained it has been processed and reconciled.
  • Lost deletes. A delete is filtered out, mishandled as an empty update, or never applied because the tombstone was dropped. The target keeps rows the source no longer has, forever.
  • Gaps. Kafka retention expires before a lagging consumer reads the events, or a connector restarts from the wrong offset. Data is missing with no error.
  • Snapshot and stream overlap. The initial snapshot and the streaming phase overlap or leave a gap at the hand-over.
  1. 01Order per key
  2. 02Deduplicate
  3. 03Apply with MERGE
  4. 04Check watermarks
  5. 05Reconcile closed windows
Figure 1. The checks in the order a batch passes through them.

Ordering: sort by the log, not the clock

The single most important design decision is what defines “newer”. Event timestamps such as ts_ms come from clocks and can tie or go backwards. The source log position, for example the PostgreSQL LSN that Debezium exposes as source.lsn, is the order the database actually committed in. I order by the log position and use the timestamp only for reporting.

Two checks follow from that.

  1. Partitioning check. All events for a primary key must be in one Kafka partition. With Debezium this is the default when the message key is the table’s primary key. Verify it rather than assume it: in a test topic, group by key and assert one distinct partition per key.
  2. Monotonicity check. Within each key in a batch, after sorting by log position, operations must form a legal sequence. A create after a create without a delete in between, or an update to a key the target has never seen and that has no create in the batch, deserves a look.
monotonic_check.sql
SELECT order_id, count(*) AS bad_transitions
FROM (
SELECT order_id, op,
lag(op) OVER (PARTITION BY order_id ORDER BY source_lsn) AS prev_op
FROM lake.staging.orders_cdc_batch
)
WHERE (op = 'c' AND prev_op IN ('c', 'u'))
OR (op IN ('u', 'd') AND prev_op = 'd')
GROUP BY order_id;

Some sequences are legitimate in particular configurations, for example a re-insert after a delete with the same key. Treat this check as a high-severity investigation trigger, not an automatic block, until you know your source’s patterns.

Deduplicate, then apply idempotently

Each micro-batch is reduced to one final event per key before it touches the target. That makes the apply step independent of how many times an event was delivered.

apply_batch.py
from pyspark.sql import Window, functions as F
def apply_batch(batch_df, batch_id):
events = (batch_df
.where(F.col("value").isNotNull()) # drop tombstones; the 'd' event carries the delete
.select(
F.coalesce(F.col("after.order_id"), F.col("before.order_id")).alias("order_id"),
F.col("op"),
F.col("source.lsn").alias("source_lsn"),
F.col("after.customer_id").alias("customer_id"),
F.col("after.amount").alias("amount"),
F.col("after.status").alias("status")))
w = Window.partitionBy("order_id").orderBy(F.col("source_lsn").desc())
latest = (events
.withColumn("rn", F.row_number().over(w))
.where("rn = 1")
.drop("rn"))
latest.createOrReplaceTempView("cdc_latest")
batch_df.sparkSession.sql("""
MERGE INTO lake.sales.orders t
USING cdc_latest s
ON t.order_id = s.order_id
WHEN MATCHED AND s.source_lsn > t.source_lsn AND s.op = 'd' THEN DELETE
WHEN MATCHED AND s.source_lsn > t.source_lsn THEN UPDATE SET
customer_id = s.customer_id, amount = s.amount, status = s.status,
source_lsn = s.source_lsn
WHEN NOT MATCHED AND s.op <> 'd' THEN INSERT
(order_id, customer_id, amount, status, source_lsn)
VALUES (s.order_id, s.customer_id, s.amount, s.status, s.source_lsn)
""")

The column parsing is simplified; in practice the Kafka value is deserialised against the registered schema first. Three details are what make it testable.

  • The target stores source_lsn. Without it, the apply step cannot tell a replayed old event from a new one.
  • Stale events are no-ops. Both WHEN MATCHED clauses require the incoming position to be newer than the one the target holds. A replayed or out-of-order event matches neither clause, so the row is left unchanged.
  • Deletes are applied, not skipped. A delete for a key not in the target is ignored rather than inserted.

On tombstones: Debezium can emit a tombstone (a record with a null value) after each delete so that Kafka log compaction can eventually remove the key. The delete itself is the preceding event with op = 'd'. The pipeline should ignore tombstones for applying changes, but a test should confirm that each delete event is present. If a downstream step filters “null” values too aggressively, it can drop both.

Test: replay is a no-op

The idempotency test is simple and catches a surprising number of bugs. Apply a batch, record the table’s checksum, apply the same batch again, and compare.

replay_idempotency.sql
SELECT checksum(concat_ws('|',
CAST(order_id AS varchar),
coalesce(CAST(amount AS varchar), '~'),
coalesce(status, '~'),
CAST(source_lsn AS varchar))) AS h,
count(*) AS n
FROM lake.sales.orders FOR VERSION AS OF 8734512094;

Run it against the snapshot after the first apply and after the replay. Both h and n must be identical. Then shuffle the batch’s event order and apply it to a fresh copy; the result must still match. If it does not, the deduplication or the stale-event guard is wrong.

Watermarks and lag

A watermark records how far the target has caught up. I keep two.

  • Source watermark. The highest log position applied to the target. Compared with the source’s current position, it gives replication lag in log terms, which is more honest than wall-clock lag during quiet periods.
  • Event-time watermark. The newest business timestamp that is complete. This is what consumers care about: “orders up to 09:00 are final”.

Checks on watermarks:

watermark_checks.sql
-- 1. The applied position must never go backwards between loads.
SELECT load_id, applied_lsn, prev_lsn
FROM (
SELECT load_id, applied_lsn,
lag(applied_lsn) OVER (ORDER BY committed_at) AS prev_lsn
FROM dq.cdc_watermarks
WHERE table_name = 'sales.orders'
)
WHERE applied_lsn < prev_lsn;
-- 2. No target row may claim a position beyond the recorded watermark.
SELECT count(*) AS rows_ahead_of_watermark
FROM lake.sales.orders
WHERE source_lsn > (SELECT max(applied_lsn) FROM dq.cdc_watermarks
WHERE table_name = 'sales.orders');

The first catches a consumer that was reset to an older offset. The second catches a watermark that is updated before the data commits, which would let a failed load declare success.

For gaps, compare Kafka offsets: the offset ranges consecutive batches committed should join up per partition, with no range skipped (individual offsets need not be contiguous, because transaction markers and compaction leave holes), and the consumer’s position should never fall behind the topic’s earliest retained offset. If it does, events were lost to retention and only a re-snapshot can recover them. That condition should be critical.

Reconciliation windows

You cannot compare a live source with a CDC target at an arbitrary moment; they differ by whatever is in flight. Reconciliation needs a window that is closed on both sides.

I define a window as a range of business time plus a settling period that covers expected lateness. For an orders table, “all orders created on 14 September, reconciled once the event-time watermark has passed 15 September 02:00” is a typical shape. The settling period is chosen from observed lateness, with margin; it is a trade-off between how early you can declare a day final and how often late events reopen it.

Within a closed window, the checks from the source-to-target reconciliation case study apply: counts, checksums and keyed diffs per partition. CDC adds one more, the operation balance.

op_balance.sql
WITH ev AS (
SELECT count_if(op IN ('c', 'r')) AS creates,
count_if(op = 'd') AS deletes
FROM lake.staging.orders_cdc_log
WHERE CAST(event_created_at AS date) = DATE '2026-09-14'
),
tgt AS (
SELECT count(*) AS live_rows
FROM lake.sales.orders
WHERE order_date = DATE '2026-09-14'
)
SELECT ev.creates, ev.deletes, tgt.live_rows,
ev.creates - ev.deletes - tgt.live_rows AS imbalance
FROM ev CROSS JOIN tgt;

For a table where each key is created and deleted at most once, creates minus deletes should equal live rows for the window. A positive imbalance points to lost creates or duplicate deletes. A negative one points to duplicate creates or lost deletes. It is a coarse check, and keys that are deleted and re-created break its assumption, but it often locates the problem class before a keyed diff runs.

For deletes specifically, an anti-join against the source catches the case that every other check misses:

ghost_rows.sql
SELECT t.order_id
FROM lake.sales.orders t
LEFT JOIN pg_sales.public.orders s ON s.order_id = t.order_id
WHERE s.order_id IS NULL
AND t.order_date = DATE '2026-09-14';

Every result is a ghost row: present in the target, gone from the source.

Test data strategies

Production CDC streams rarely contain the hard cases on demand. I build them.

  • Scripted scenario tables. A small source database where a test harness executes known sequences per key: insert, update, update, delete, re-insert. Each scenario has an expected final target state.
  • Fault injection at the stream. Replay a batch, drop a partition’s events, swap the order of two events for a key, delay a subset beyond the settling period. Each fault must be caught by a named check.
  • Snapshot hand-over tests. Start the connector while the scenario script is writing, so that the snapshot and streaming phases overlap. The target must still match the source after the window closes.
  • Schema change mid-stream. Add a nullable column and then retype one while events flow. The pipeline should either evolve the target according to policy or stop with a clear error, not write nulls silently.
  • Volume tests. Generate high-churn keys, where the same key is updated thousands of times per batch, to confirm deduplication is correct and does not explode the shuffle.

These scenarios become a regression suite that runs whenever the apply logic, connector configuration or table format changes. They are also where the rule catalogue described in Part 4 gets its evidence: a rule earns its severity by catching a fault it was designed for.

Where I am still refining

Two areas remain open in my own thinking. First, settling periods are usually chosen once and forgotten; I would like them derived from measured lateness and revised automatically. Second, the operation-balance check is weak for tables with heavy delete-and-recreate patterns. A per-key state comparison at window close is more precise but costs a keyed diff every time. I treat both as trade-offs to revisit per table rather than settled answers.

Further reading

Big data testing & data quality

Try “evaluation”, “red-teaming”, “governance” or “agents”.