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

Previously

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

  1. Watch
  2. Try it
  3. Predict
  4. Capture
Debezium c…tails WAL ·…topic: ledger.eventsat-least-once deliveryPARTITIONS · 1P0→ ledger groupCONSUMER GROUP · LEDGERtwo sinks reading the same partitionSink A · UPSERTidempotent · dedup by eventId · balance = 0Sink B · balance += amountnon-idempotent · balance = 0
What to watch for

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.

Continue unlocks when the animation finishes.
Implementation

Highlighted lines are the ones running in the diagram right now.

SinkA.applyEvent # idempotent
UPSERT keyed on a stable dedup key — replays no-op.
def apply_event(event):
# dedup key chosen at sink config time
if dedup_mode == 'eventId':
key = event.eventId # outbox: SMT-assigned UUID
else:
key = (event.table, event.pk, event.lsn) # raw CDC
db.execute('''
INSERT INTO accounts(account_id, balance)
VALUES (%s, %s)
ON CONFLICT (account_id) DO UPDATE
SET balance = EXCLUDED.balance
WHERE accounts.dedup_key IS DISTINCT FROM %s
''', [event.pk, event.after_balance, key])
SinkB.applyEvent # non-idempotent
balance += amount — every delivery moves state.
def apply_event(event):
# no dedup key, no UPSERT, no fingerprint check
db.execute('''
UPDATE accounts
SET balance = balance + %s
WHERE account_id = %s
''', [event.amount, event.pk])
# replay drifts balance by replayed sum
enableEosOnDebezium # the misleading switch
EOS scope ends at Kafka — does not extend to the sink.
def configure_connector():
# Kafka 3.3+ source-connector EOS
connector.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

Built with Arqly
Every scene in Build a CDC pipeline (Debezium + outbox) builds on the one before it.All 12 Build a CDC pipeline (Debezium + outbox) scenes