Skip to content
AITroveRead. Build. Understand.
Make this comfortable

Change data capture: apply updates and deletes by source order

Last updated: 6 Oct 20265 min read
tutorial
IntermediateBy AITrove Editorial

A CDC consumer applies insert, update and delete events with source ordering and replay protection, then advances its offset only after the target state is durable.

Treat a delete as data

A receipt correction stream may emit a delete when a record is withdrawn. If the consumer ignores delete events, downstream tables keep a receipt that no longer exists at the source. A tombstone can remove the row or mark it inactive according to the warehouse retention policy. The business key and source position are part of every change; arrival order over the network is not necessarily commit order.

Guard against replay and reordering

Keep the last applied source position per key or per ordered partition. A duplicate event should have no second effect; an older update must not revive a row after a newer delete. If the source splits one key across partitions without a global order, establish a consistent partitioning or version rule first. Target and checkpoint ownership] must be transactional or reconciled after failures.

Bootstrap consistently

An initial snapshot and a live change stream need a handoff position. Start the stream from a point consistent with the snapshot; otherwise changes during copying can be missed or applied twice. Applying twice is safe only if the target operation is idempotent. Reconcile row counts and sampled keys after bootstrap, including a record updated and then deleted during the handoff.

Make history policy explicit

A raw change log may retain an audit trail while a current-state table exposes only the latest row. Do not delete raw evidence to satisfy a curated current-state query, and do not expose deleted sensitive content beyond its retention policy. Historical dimensions] use ordered changes for a different purpose: reconstructing a past view.

Implementation

python
def apply_change(current, event):
    prior = current.get(event.receipt_id)
    if prior is not None and event.source_position <= prior["position"]:
        return False
    current[event.receipt_id] = {
        "position": event.source_position,
        "deleted": event.operation == "DELETE",
        "value": None if event.operation == "DELETE" else event.payload,
    }
    return True

Performance and operating cost

Applying E events is O(E) expected time with keyed state; durable position state grows with active or tombstoned keys. Retention and compaction need a rule that cannot allow an old replay to resurrect forgotten keys.

Common Mistakes

  • Do not ignore deletes because the target is analytical.
  • Do not use consumer arrival time as source commit order.
  • Do not start a change stream after a snapshot without a consistent handoff position.

Read next

Continue the workflow: Idempotent streaming sinks: make retries safe across failures.

Continue the workflow: Knowledge-graph provenance: claim time, evidence and retraction.

Continue the workflow: Deletion propagation and erasure proof.

Continue the workflow: Event envelopes, identity and tombstones.

ai-data
data-engineering
Storage details