Overview
CDC keeps distributed data in sync in near-real-time. Reconciliation proves that it worked — and catches the cases where it didn’t. Every production CDC deployment eventually produces a break: a row that was updated on the source while the Debezium connector was down, a slot replay that delivered an event twice, an outbox row deleted before the connector picked it up because someone ran a manual cleanup. These breaks are rare and usually small. In a regulated bank they are also mandatory to detect, investigate, and report within a prescribed window.
SAMA’s Technology Risk Management (TRM) framework does not prescribe a reconciliation technology stack. It does require that banks demonstrate data integrity controls across system boundaries — evidence that data transmitted between core systems is complete, accurate, and that discrepancies are detected and resolved. A Debezium connector that “mostly works” is not a control; a reconciliation job with a break report and an investigation log is.
This article covers the full pipeline: fingerprint generation, hash comparison at scale, intraday vs. daily scheduling, automated break classification, and the reporting format that satisfies an examination team.
Connector lag, offset freshness, and replication slot disk usage are operational signals that tell you the pipeline is running. Reconciliation is a data signal that tells you the pipeline produced the right result. Both are required; neither substitutes for the other.
CDC is not an audit story
CDC captures the log of what the database did. Reconciliation verifies the outcome in the downstream. The gap between them is where the breaks live.
| Signal | What it tells you | What it misses |
|---|---|---|
| Connector lag (LSN gap) | Events are flowing / backed up | Whether downstream applied them correctly |
| Consumer offset lag | Consumer is keeping up | Whether processed events produced correct state |
| Consumer error rate | Processing is failing | Events silently skipped due to bad idempotency logic |
| Slot reset event | Replay started | Whether replay covered all rows correctly |
| Reconciliation hash mismatch | Downstream diverged from source | Nothing — this is the ground-truth check |
Four break categories arise in practice:
- Missing row. Source has a record the downstream does not. CDC event was lost or dropped.
- Stale row. Source has a newer version of a record than the downstream. An update event was skipped or arrived out of order.
- Phantom row. Downstream has a record the source no longer has. A delete event was missed, or an insert was duplicated without deduplication.
- Field divergence. Both systems have the row; one or more fields differ. Transformation bug or partial update applied twice.
Comparing two tables by hand each morning with a shared spreadsheet is a control that can’t be audited, can’t scale to millions of rows, and produces no timestamp evidence. SAMA expects an automated, repeatable process with a stored result log. Manual checks are evidence that the control does not exist.
Pipeline architecture
A production reconciliation pipeline has two independent paths: a near-real-time streaming path for intraday alerting, and a batch path for the authoritative daily report. The batch path is what the regulator examines; the streaming path is what the operations team monitors.
Fingerprinting & hash comparison
The core technique is row fingerprinting: for each row in the source and downstream, compute a deterministic hash of the business-significant columns; compare hashes by primary key. A mismatch on the hash means a break; matching hashes mean the row is in sync. At SAIB’s scale (tens of millions of payment records), hashing allows the comparison to run in minutes rather than column-by-column SQL joins that would take hours.
Do not hash technical columns like updated_at, version, or _kafka_offset — these legitimately differ between source and downstream by design. Hash only columns that must be identical: amounts, account numbers, reference IDs, status codes, ISO 20022 field values. Document the exclusion list explicitly; it is part of the reconciliation control definition.
-- Source DB fingerprint: one hash per payment row
SELECT
payment_id,
md5(
CONCAT(
COALESCE(amount::text, ''), '|',
COALESCE(currency, ''), '|',
COALESCE(debtor_iban, ''), '|',
COALESCE(creditor_iban, ''), '|',
COALESCE(status, ''), '|',
COALESCE(value_date::text, '')
)
) AS row_hash,
updated_at -- used to set the reconciliation window boundary
FROM payments
WHERE updated_at < :window_end -- stable window: no rows still in-flight
AND updated_at >= :window_start; -- incremental mode: only changed rows
-- Downstream fingerprint (Elasticsearch via JDBC bridge or export)
-- Same query shape; same column set; same COALESCE defaults
The comparison itself is a hash join on payment_id. Three result sets emerge:
- Source only — missing rows in downstream
- Downstream only — phantom rows in downstream
- Hash mismatch — field divergence
// KeyedCoProcessFunction comparing source and downstream hash streams
public class HashCompareFunction
extends KeyedCoProcessFunction<String, Fingerprint, Fingerprint, Break> {
// State: hold source hash until downstream arrives (or TTL fires)
private ValueState<Fingerprint> sourceState;
private ValueState<Fingerprint> downstreamState;
@Override
public void processElement1(Fingerprint src, Context ctx, Collector<Break> out)
throws Exception {
Fingerprint ds = downstreamState.value();
if (ds == null) {
sourceState.update(src);
// register a timer: if no downstream arrives, emit MISSING break
ctx.timerService().registerEventTimeTimer(src.windowEnd + 60_000);
} else if (!src.hash.equals(ds.hash)) {
out.collect(Break.fieldDivergence(src, ds));
clearState();
} else {
clearState(); // hashes match, no break
}
}
@Override
public void onTimer(long timestamp, OnTimerContext ctx, Collector<Break> out)
throws Exception {
Fingerprint src = sourceState.value();
if (src != null) out.collect(Break.missing(src));
clearState();
}
}
Scheduling & windowing
The hardest design question in reconciliation is when to cut the window. Cut it too early and in-flight events (written to source, not yet propagated via CDC) look like breaks. Cut it too late and the report is not ready for the regulatory deadline.
The window boundary must be set at least as far back as the measured p99 end-to-end CDC latency — not a round number. If your Debezium-to-downstream p99 is 8 seconds at peak, a window that cuts 5 seconds before query time will consistently flag in-flight events as breaks. Use a 30-second or 60-second lag behind the current wall clock, validated against actual measurements.
The two scheduling strategies:
- Intraday (streaming): 15-minute tumbling windows, Flink watermark set at window-end minus the measured propagation lag. Result: a near-real-time break count, alerted if it exceeds a threshold (e.g., more than 0 unexplained breaks in a 15-minute window). This feeds the operations dashboard.
- Daily batch (authoritative): Runs at 02:00 AST on the prior business day’s data. The window is T−1 business day, 00:00:00–23:59:59 AST, with a 90-second lag. Spark reads source and downstream via JDBC. Result is written to the
recon_runstable and signed with a SHA-256 of the result set for tamper evidence. This is the regulatory artefact.
Automated break investigation
A break list without a root cause is not an actionable artefact. Automated investigation classifies each break before it reaches the operations team, eliminating the category of breaks that have known-safe explanations.
-
Check if the break is in a CDC propagation window
If the source row’s
updated_atis within the propagation lag window of the reconciliation cut-off, the break may be an in-flight event. Flag asPENDING_PROPAGATIONand re-evaluate in the next run rather than alerting immediately. -
Correlate against the connector restart log
If the Debezium connector restarted or a slot was reset in the reconciliation window, breaks are expected during the replay period. Tag affected rows as
CONNECTOR_RESTARTand mark as auto-investigated. The connector restart is itself an incident; the breaks are a downstream symptom. -
Check for known schema migration windows
A running Flyway migration or Liquibase changeset may produce temporary divergence as the downstream schema catches up. If a deployment occurred during the window, mark affected rows
SCHEMA_MIGRATION. Confirm the downstream has converged in the next run. -
Re-query the source for the current value
For a row flagged as missing or stale, query the source directly. If the row no longer exists (deleted after the reconciliation window closed), reclassify as
POST_WINDOW_DELETE— not a break. If the row exists and differs from the downstream, escalate toOPEN_BREAK. -
Attempt self-heal via replay
For missing and stale rows that remain unexplained, publish a replay request to the
recon.replay-requestsKafka topic. The downstream consumer re-fetches the row from the source API and applies it. Log the replay attempt with timestamp and outcome. Do not retry more than once without human sign-off; infinite replays mask a systematic fault. -
Escalate OPEN_BREAK items
Any break that survives all the above steps is an
OPEN_BREAKrequiring human investigation. Write to thebreak_investigationstable, fire a PagerDuty alert, and set an SLA timer. SAMA TRM requires that data integrity failures are investigated and resolved within a documented timeframe — the timer enforces it.
SAMA audit trail requirements
SAMA’s Technology Risk Management framework and the Cyber Security Framework both treat data integrity as an auditable control. In practice, an examination team will ask for evidence of three things:
- That reconciliation runs on a defined schedule. The
recon_runstable with itsscheduled_at,started_at,completed_at, andstatuscolumns is this evidence. A run that was skipped must have a reason recorded. - That breaks are detected, classified, and resolved within the SLA. The
break_investigationstable with itsdetected_at,classified_at,resolved_at,resolution_type, andresolved_bycolumns provides the case log. - That the reconciliation result cannot be altered after the fact. The SHA-256 hash of the result set, stored in
recon_runs.result_hash, provides tamper evidence. Store the raw result in object storage (S3 or Ceph) under an immutable retention policy.
CREATE TABLE recon_runs (
run_id uuid PRIMARY KEY DEFAULT gen_random_uuid(),
pipeline_name varchar(128) NOT NULL,
window_start timestamptz NOT NULL,
window_end timestamptz NOT NULL,
scheduled_at timestamptz NOT NULL,
started_at timestamptz,
completed_at timestamptz,
status varchar(32) NOT NULL DEFAULT 'PENDING',
source_row_count bigint,
downstream_row_count bigint,
break_count int NOT NULL DEFAULT 0,
result_hash varchar(64), -- SHA-256 of result CSV
result_object text, -- s3://recon-archive/runs/{run_id}.csv.gz
skip_reason text
);
CREATE TABLE break_investigations (
break_id uuid PRIMARY KEY DEFAULT gen_random_uuid(),
run_id uuid REFERENCES recon_runs,
entity_id varchar(128) NOT NULL,
break_type varchar(32) NOT NULL, -- MISSING | PHANTOM | STALE | FIELD_DIVERGENCE
detected_at timestamptz NOT NULL DEFAULT now(),
auto_classified boolean NOT NULL DEFAULT false,
classification varchar(64), -- PENDING_PROPAGATION | CONNECTOR_RESTART | OPEN_BREAK …
replay_attempted boolean NOT NULL DEFAULT false,
resolved_at timestamptz,
resolution_type varchar(64),
resolved_by varchar(128),
sla_deadline timestamptz GENERATED ALWAYS AS
(detected_at + INTERVAL '4 hours') STORED
);
SAMA TRM requires data integrity records to be retained for a minimum of 7 years. The recon_runs and break_investigations tables must be on a retention policy that prevents deletion or modification by application code. Use PostgreSQL row-level security to block DELETEs from the reconciliation service role, and store the result CSVs in an S3 bucket with Object Lock in Compliance mode.
Building the pipeline
-
Define the reconciliation domain and column set
Write a reconciliation specification document that names: the source table, the downstream store, the primary key, the columns included in the hash (and excluded), the propagation lag SLA, the window schedule, and the break SLA. This document is a control artefact; version it in git and review it with Risk and Compliance before go-live.
-
Build the fingerprint extractors
Source extractor: a Spark job (for batch) or a Flink source reading from the CDC topic alongside a periodic JDBC snapshot for comparison. Downstream extractor: same column set, same hash function, same
COALESCEdefaults. Run both extractors against a synthetic dataset and assert that identical rows produce identical hashes — this is the most common source of false-positive breaks. -
Stand up the break store schema
Apply the
recon_runsandbreak_investigationsDDL via a versioned migration. Create the application role with INSERT/UPDATE on both tables and SELECT on everything, but no DELETE and no TRUNCATE. Apply row-level security for the mutation restriction. -
Implement and test the break investigation logic
The classification pipeline is a state machine. Test each transition with known break scenarios from a staging environment. Run a chaos test: deliberately stop the Debezium connector for 5 minutes, restart, let the replay run, then verify that reconciliation classifies the affected rows as
CONNECTOR_RESTARTand not asOPEN_BREAK. -
Wire the alerting and SLA monitoring
Alert on: any
OPEN_BREAKcreated; any break that reaches 50% of the SLA deadline without aresolved_at; any run that does not start within 15 minutes of its scheduled time. All three are SLA violations under SAMA TRM; all three should page the on-call engineer, not just create a Jira ticket. -
Generate and store the result artefact
At the end of each batch run: serialize the full break list as a gzip-compressed CSV; compute the SHA-256 of the raw bytes; upload to S3 with server-side encryption (SSE-S3 at minimum, SSE-KMS preferred); write the object key and hash to
recon_runs.result_objectandrecon_runs.result_hash. The upload must complete before the run is markedCOMPLETED. -
Run the first production reconciliation under observation
Do not run the self-heal replay step in production until the break classification accuracy has been validated — an over-eager replay that applies stale source data to a downstream that is actually correct creates a new break rather than resolving one. Disable replay for the first two weeks; classify and manually resolve breaks; then enable replay only for the
MISSINGbreak type.
Regulator-ready reporting
The daily report is a two-page summary: a run statistics block and a break detail section. It is produced from the break store and emailed to a distribution list that includes Risk and Compliance. A zero-break day is still a report — the absence of breaks is itself a data point and the evidence that the pipeline ran.
-- Summary for the prior business day report
SELECT
r.pipeline_name,
r.window_start AS window_start,
r.window_end AS window_end,
r.source_row_count,
r.downstream_row_count,
r.source_row_count - r.downstream_row_count
AS count_variance,
r.break_count,
COUNT(b.break_id) FILTER (WHERE b.classification = 'OPEN_BREAK')
AS open_breaks,
COUNT(b.break_id) FILTER (WHERE b.resolved_at IS NOT NULL)
AS resolved_breaks,
COUNT(b.break_id) FILTER
(WHERE b.sla_deadline < now() AND b.resolved_at IS NULL)
AS sla_breached,
r.result_hash,
r.completed_at
FROM recon_runs r
LEFT JOIN break_investigations b ON b.run_id = r.run_id
WHERE r.window_start::date = current_date - 1
GROUP BY r.run_id
ORDER BY r.pipeline_name;
Include the result_hash in the email body. If an examiner queries the break list at a later date, they can verify that the report was not modified after the run completed by recomputing the hash from the archived CSV. This is the chain of custody that turns a database query into an auditable control.
Common pitfalls
A NULL in the source that appears as an empty string in the downstream (due to ORM serialization or JSON coercion) produces a hash mismatch on every row that has any NULL in the hashed column set. Always normalize NULLs to a sentinel value (e.g., empty string) using COALESCE in both extractors. Validate this with a dataset that contains known NULLs before go-live.
If the Debezium connector replays events (slot reset, snapshot), the downstream is in a partially-applied state during the replay. A reconciliation run over that window will produce a large number of false breaks. Implement a reconciliation skip mechanism that detects an in-progress replay (check the connector’s snapshot mode from the Kafka Connect REST API) and records SKIPPED_DURING_REPLAY in recon_runs.skip_reason. Resume the next scheduled run after replay completion.
SAMA reporting deadlines are in AST (UTC+3). If the source DB uses UTC and the downstream uses a different timezone, a naive window boundary will include or exclude the wrong rows at day boundaries. Store all timestamps in UTC in the break store; convert to AST only in the report output layer. A one-hour error in the window boundary can cause an entire hour’s worth of rows to appear as breaks.