At-least-once is the ceiling
Debezium delivers at-least-once and Kafka EOS only covers Debezium↔Kafka, so end-to-end correctness depends on the sink being idempotent — typically by deduplicating on a stable key like eventId or (table, pk, lsn).
Within-key ordering is preserved end-to-end — but Debezium's delivery guarantee is at-least-once, so the consumer will see the same ordered event more than once on retries and restarts unless it is built to deduplicate.
Scene 11
At-least-once is the ceiling
- Watch
- Try it
- Predict
- Capture
Debezium tails the ledger table's WAL and emits one change event per row change. Both sinks agree on the running balance — they're seeing each event exactly once, in order. The interesting question (next phase) is what happens on a connector restart.
Highlighted lines are the ones running in the diagram right now.
def apply_event(event):# dedup key chosen at sink config timeif dedup_mode == 'eventId':key = event.eventId # outbox: SMT-assigned UUIDelse:key = (event.table, event.pk, event.lsn) # raw CDCdb.execute('''INSERT INTO accounts(account_id, balance)VALUES (%s, %s)ON CONFLICT (account_id) DO UPDATESET balance = EXCLUDED.balanceWHERE accounts.dedup_key IS DISTINCT FROM %s''', [event.pk, event.after_balance, key])
def apply_event(event):# no dedup key, no UPSERT, no fingerprint checkdb.execute('''UPDATE accountsSET balance = balance + %sWHERE account_id = %s''', [event.amount, event.pk])# replay drifts balance by replayed sum
def configure_connector():# Kafka 3.3+ source-connector EOSconnector.set('exactly.once.support', 'required')# Fences Debezium ↔ Kafka writes inside a transaction.# NOT covered by this switch:# - the sink's reads from Kafka# - the sink's writes to its own store# The sink is outside the transaction; it must dedup itself.
Where this sits in Build a CDC pipeline (Debezium + outbox)
Scene 11 of 12. Debezium delivers at-least-once and Kafka EOS only covers Debezium↔Kafka — end-to-end correctness depends on the sink being idempotent, typically by deduplicating on (table, pk, lsn) or eventId.
Up next. Every piece is on the table — log, slot, snapshot, schema registry, outbox, cleanup, partitioning, idempotency — and the next move is to compose them into a pipeline for a real workload and trace each choice back to the failure it prevents.
All 12 scenes in Build a CDC pipeline (Debezium + outbox) · Every curriculum