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.

This complements, not replaces, CDC monitoring

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.
Manual SQL queries are not a reconciliation pipeline

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.

Which columns to hash

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.

recon-fingerprint.sqlsql
-- 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
ReconJob.java (Flink streaming)java
// 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.

CDC propagation lag must be measured, not assumed

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_runs table 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.

  1. Check if the break is in a CDC propagation window

    If the source row’s updated_at is within the propagation lag window of the reconciliation cut-off, the break may be an in-flight event. Flag as PENDING_PROPAGATION and re-evaluate in the next run rather than alerting immediately.

  2. 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_RESTART and mark as auto-investigated. The connector restart is itself an incident; the breaks are a downstream symptom.

  3. 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.

  4. 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 to OPEN_BREAK.

  5. Attempt self-heal via replay

    For missing and stale rows that remain unexplained, publish a replay request to the recon.replay-requests Kafka 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.

  6. Escalate OPEN_BREAK items

    Any break that survives all the above steps is an OPEN_BREAK requiring human investigation. Write to the break_investigations table, 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_runs table with its scheduled_at, started_at, completed_at, and status columns 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_investigations table with its detected_at, classified_at, resolved_at, resolution_type, and resolved_by columns 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.
recon-schema.sqlsql
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
);
Retention: 7 years, immutable

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

  1. 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.

  2. 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 COALESCE defaults. 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.

  3. Stand up the break store schema

    Apply the recon_runs and break_investigations DDL 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.

  4. 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_RESTART and not as OPEN_BREAK.

  5. Wire the alerting and SLA monitoring

    Alert on: any OPEN_BREAK created; any break that reaches 50% of the SLA deadline without a resolved_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.

  6. 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_object and recon_runs.result_hash. The upload must complete before the run is marked COMPLETED.

  7. 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 MISSING break 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.

daily-recon-report.sqlsql
-- 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

Hashing NULL differently on source and downstream

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.

Running reconciliation during a CDC replay

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.

Timezone assumptions in window boundaries

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.