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.
Change data capture: apply updates and deletes by source order
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
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 TruePerformance 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
- Idempotent loads: commit target rows and extraction progress together
- Incremental extraction: advance a compound watermark without losing ties
- Warehouse history: join facts to the dimension version valid at event time
- Data quality gates: quarantine bad rows and reconcile complete batches
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.
