Continuously moving database changes downstream delivers only part of a data product. A report must also explain which business period a number represents, whether it includes refunds, how late records revise it, and whether recovery changes the result. This article uses a hypothetical transaction system to explore the design and validation required to turn a CDC pipeline into a trusted product. All example data and thresholds are illustrative.
1. Define the business facts before moving the changes
Imagine order O42 receives a payment of CNY 100 and later a refund of CNY 20. The orders table holds current state, the payments table records receipts, and the refunds table records refund completion times. The business needs two products: customer support needs the order’s current net receipts of CNY 80, while finance needs daily cash movements attributed to the dates payments and refunds occurred. They share source tables but differ in grain, time, and revision rules. A single wide order table cannot answer every question.
CDC captures changes to stored records. A status update might represent a business action, a manual correction, a historical migration, or the outcome of a duplicate callback. Counting every update as a new sale confuses technical activity with business facts. Start by specifying the product key, monetary units, state definitions, time fields, and correction rules. Then establish whether the available changes contain enough information to support them.
A current-state product can use the order ID as its unique key and maintain the latest valid version. A cash-movement product is better served by payment and refund records with stable identifiers. If the source overwrites those records and provides no reconstructable history, the product cannot promise to reproduce every past state. Add a business ledger or narrow the promise when that boundary appears; a more elaborate consumer cannot recover information that never existed.
2. Give snapshots and incremental changes a defensible boundary
On initial ingestion, historical rows come from a snapshot and new changes come from the log. Together they must form a complete state. Debezium’s PostgreSQL connector documentation describes how snapshots connect to log positions. Acceptance still requires checking the snapshot mode, isolation settings, and recovery behavior actually in use. An arbitrary full-table export combined with incremental consumption starting at an arbitrary time does not demonstrate that nothing was missed.
Suppose the snapshot reads O42 with net receipts of CNY 100 while a refund changes the value to CNY 80 during the scan. The final state must be CNY 80, and reading the snapshot row must not create CNY 100 of new sales. Snapshot records describe existing state at a consistency boundary; incremental records describe subsequent changes. Preserve provenance and snapshot batch identifiers, and define how overlapping records are merged.
Rehearse failure halfway through the scan. Recovery may reread rows that were already scanned, so consumers must tolerate duplicates. Retention for snapshots, source logs, and consumer positions must cover the real recovery window. If required logs are gone, enter an explicit rebuild procedure. One option is to populate a separate version, verify completeness, and then switch serving to it, rather than exposing a mixture of half a new snapshot and half an old state.
3. Separate change ordering from business time
If an order’s paid version V12 and refunded version V13 arrive in reverse order because of retries or different processing paths, simple last-arrival-wins behavior restores an obsolete paid state. State comparisons need a reliable source ordering that distinguishes events within the same transaction. A source identity or failover epoch may also be necessary. Receipt time, application update time, and message partition positions cannot be treated interchangeably as a global version.
Define the scope of every comparison. A position orders records within its partition; it does not automatically compare with a position in another partition. Log positions from separate databases also have no inherent ordering. For a business key, establish a single authoritative writer or an explicit conflict rule. After primary failover or a backfill, verify that versions remain comparable. Quarantine incomparable records instead of silently overwriting state.
Business time addresses a different problem: a refund occurring on September 26 may not reach analytics until September 27. Flink documents watermarks as a measure of event-time progress while acknowledging that records can arrive late. A daily report can therefore distinguish provisional values from values after closing and define whether late records trigger recomputation or a recorded revision. Watermarks control waiting; they cannot prove that a source will never enter another historical transaction.
4. Deletion needs business meaning and a retention policy
Deleting an order row does not automatically mean a refund. It may remove a test order, archive old data, or fulfill a cleanup requirement. If a consumer subtracts revenue for every deletion, an archive job could erase historical performance overnight. The product contract must distinguish business cancellation, soft deletion, physical deletion, and historical retention, and specify which affect current queries or historical metrics.
Debezium’s PostgreSQL documentation distinguishes delete events from null-valued tombstones used for log compaction. Verify which representation reaches the consumer and whether intermediate transformations remove keys or deletion indicators. Kafka’s design documentation also explains that observing delete markers depends on retention and consumer progress. A long-paused consumer resuming successfully does not, by itself, prove that all previously deleted rows have been removed from its old state.
A current-state store can retain a deletion version to prevent a late, older update from resurrecting a row. Keep that version for a period consistent with the history that replay permits. If personal data must be fully removed, also examine copies in detailed datasets, indexes, caches, and recoverable history. Invisibility in business queries and removal from physical copies are separate acceptance criteria and need separate evidence. A single deletion message cannot establish both.
5. Extend idempotency to final writes and derived calculations
A familiar failure window occurs after the target write succeeds but before the consumer position is persisted. The process exits, and recovery delivers the same event again. Replacing current state by primary key may remain correct, while adding its amount to daily sales counts it twice. Idempotency must cover the final side effect. A delivery guarantee in the messaging system does not automatically protect an external database, cache, or notification.
One approach commits a deduplication record and the business write within the same target transaction, using a stable event identifier that distinguishes source changes. Another tracks each order’s applied version and previous contribution. When O42 changes from CNY 100 to CNY 80, retract the old contribution and apply the new one; reapplying the same version changes nothing. If a grouping dimension changes too, subtract from the old group and add to the new one.
Keeping only the latest version cannot necessarily reconstruct every cash movement, because intermediate states may contain independent business facts. Separate current state, an immutable event ledger, and aggregates according to their purpose. When choosing transactional writes, conditional updates, or rebuildable aggregates, explain failure atomicity and storage costs. If an external system cannot support the required transaction boundary, keep reconcilable task state and arrange compensation instead of making an unverifiable end-to-end promise.
6. Recognize business objects that are temporarily incomplete
An order header, its items, and a payment record may originate in one transaction yet arrive through different topics and consumers. If a paid order arrives before its payment amount, an immediate join produces a missing amount. Interpreting that absence as zero creates a false revenue drop. Decide whether the product may display partial state or must wait until related records meet completeness conditions. This is a business availability decision.
Possible strategies include processing around transaction boundaries, maintaining source-table state and recomputing affected orders, or holding incomplete objects in a pending area. A 30-second waiting limit for a support page, for example, is merely a design value to test. After it expires, display that data is incomplete and retain a retry path; expiration is not evidence that no payment exists. Workflows across databases have no inherent shared commit instant and need business identifiers or explicit coordination.
Validate join cardinality as well. O42 has two items and two refunds. Joining all three tables directly and summing the order amount can count it four times. Aggregate at each table’s business grain before joining on verified keys. Dimension records without a unique match should enter an exception set. Row counts alone are unlikely to reveal every inflated amount, so reconciliation must check both uniqueness and business totals.
7. Make the quality contract describe failure behavior
The Open Data Contract Standard includes schema, data quality, and service expectations in the agreement between producers and consumers. An order product’s implementation should specify at least that order IDs are present and unique in current state, which currencies and precision amounts use, how refunds affect net receipts, how status values evolve, how freshness is measured, and who handles violations. Documentation is a starting point; rules need executable checks or an explicit verification process.
Give rules different consequences. A missing business key breaks idempotency and may justify blocking the affected record and alerting an owner. A missing optional note normally need not stop the whole product. A new source status must not silently map to completed: mark its meaning as unknown and flag affected metrics for review. Adding a nullable column and changing a monetary unit from yuan to fen may both look like schema changes, but their compatibility implications are entirely different.
Freshness cannot be measured solely by whether a consumer is active. An old last business event may be normal when the source has no transactions. New source transactions without advancing capture progress indicate a different condition. Combine source heartbeats, capture progress, processing backlog, and the serving layer’s update time. A contract can require displaying the data cutoff and reducing its trust status when delays exceed the limit; a green process-health icon is no substitute for product availability.
8. Prove delivery with fault injection, replay, and reconciliation
An acceptance dataset should include creation, consecutive updates to one key, duplicates, reordering, old updates after deletion, primary-key changes, refunds across dates, and schema changes. Specify both expected current state and expected metrics for each case. Checking that messages arrived is insufficient. Inject failures before and after target writes, around position commits, and at snapshot cutover; observe whether recovery converges on the same result.
Replay validation needs fixed input boundaries, code versions, and business definitions. Repeating the same changes, or varying arrival order within supported conditions, should preserve declared invariants. Reconciliation needs a shared cutoff or a reconstructable source snapshot; yesterday’s downstream state cannot be directly compared with a source that continues changing. Beyond total row counts, compare key sets, monetary totals, and detailed differences by date, status, and currency.
The delivered product should let a user interrogate a number: which definition version produced it, what its data cutoff is, which partitions remain incomplete, why a revision happened, and how to rebuild it. When those questions have operational answers, CDC becomes a dependable data product. Explicitly marking insufficient source information, expired history, and unresolved business meaning is itself part of that quality capability.
