Data engineering
Part 4 of 5 · CDC & DebeziumExactly-Once CDC Pipelines — Idempotent Consumers, Keys & At-Least-Once Reality
CDC capture is at-least-once. Exactly-once effects come from a stable key plus an idempotent sink: an LSN guard for projections, an inbox for side effects. Kafka transactions do not make an email or a charge exactly-once.
- 1Gist
- 2Maps
- 3Q&A
- 4Sandbox
Voice readout needs Web Speech Synthesis in this browser.
Effects, not a magic transport
Prefer
At-least-once delivery plus an idempotent sink
Snapshots, connector restarts, and consumer rebalances all redeliver. A stable key makes the second apply a no-op.
- Projections compare LSN or version and skip stale events.
- Charges and email use an inbox unique on event id.
- Ack moves only after the sink commit.
Alternative
The vendor said exactly-once
That phrase usually describes a connector offset or a Kafka transaction. It does not describe your HTTP call.
- Ignoring duplicates corrupts counts and double-charges.
- An in-memory set dies on restart and ignores other instances.
- Ack before the sink write loses the event.
Overview
Vendors say “exactly-once.” Seniors unpack the layers. The database transaction commits once. The connector’s offset flush redelivers after a crash. The consumer redelivers if it crashes before ack. The sink is exactly-once only when a second delivery does not change the outcome.
CDC makes this sharper than a hand-rolled producer. Snapshots emit op=r, then streaming may emit the same row again. Restarts replay from the last flushed offset. You design for that on purpose.
This page applies those ideas to CDC envelopes. It does not re-teach Kafka transactional producers. Read Delivery Semantics for transactional.id and read-committed. Read Transactional Outbox & Inbox Patterns for the case where the event is a domain fact co-committed with the order, and for the inbox table shape.
The honest layering
| Layer | Typical guarantee | What you control |
|---|---|---|
| DB commit | Once, or not at all | The business transaction |
| CDC capture to the bus | At-least-once (offset flush) | Connector health and flush policy |
| Bus to consumer | At-least-once, unless the sink is Kafka and you use EOS carefully | Consumer design |
| Sink effect | Exactly-once if the apply is idempotent | Keys, versions, inbox, merges |
Interview line: exactly-once effects, not magic exactly-once transport.
Stable keys
- Message key = entity primary key, or aggregate id from the outbox. Same key, same partition, compaction can drop older values. Partition arithmetic stays in Partition Keys — Ordering Guarantees vs Parallel Throughput. Do not re-derive it here.
- Tombstone: a delete for that key with a null value, so a compacted topic can forget the key. Without it, an old value can resurrect after compaction, or the sink never deletes.
- Idempotency token: pick one scheme and keep it. Outbox
event_id, or(table, primary key, LSN, op). Do not mint a new id at consume time.
Decisions
- 1
1. CDC event arrives
- next2. Stable event id or LSN
- ?
2. Stable event id or LSN
- no3. Stop and redesign identity
- yes4. What kind of sink
- 3
3. Stop and redesign identity
- ?
4. What kind of sink
- projection5. Upsert with an LSN guard
- side effect6. Inbox unique on event id
- warehouse7. Merge on key plus batch id
- 5
5. Upsert with an LSN guard
- next8. Commit sink then ack
- 6
6. Inbox unique on event id
- next8. Commit sink then ack
- 7
7. Merge on key plus batch id
- next8. Commit sink then ack
- 8
8. Commit sink then ack
Lesson map
Exactly-Once CDC Pipelines — Idempotent Consumers, Keys & At-Least-Once Reality
CDC capture is at-least-once. Exactly-once effects come from a stable key plus an idempotent sink: an LSN guard for projections, an inbox for side effects. Kafka transactions do not make an email or a charge exactly-once.
Architecture. Architecture
Select a node to see why it exists, or an edge to see the protocol, direction, effect, and consequence.
Mermaid export
flowchart TB e["1. CDC event arrives"] k["2. Stable event id or LSN"] bad["3. Stop and redesign identity"] i["4. What kind of sink"] e -->|1. CDC event arrives| k k -->|no| bad k -->|yes| i
Idempotent sink patterns
- Version or LSN guard.
UPDATE … WHERE lsn < :incoming. A redelivery of an older or equal LSN is a no-op. This fits projections that are a function of the row. - Inbox. Insert
event_idunder a unique constraint in the same transaction as the side effect. Conflict means skip. This fits email, charges, and anything that is not a pure overwrite. The table shape is on the outbox lesson. - Last-write-wins overwrite. Only if per-key order is real and you have thought about snapshot-versus-stream races. Blind overwrite of a stale redelivery is how status flaps backward.
- Warehouse merge. Stage, then
MERGEon the primary key. Store an ingestion id so a replay is auditable.
| Strategy | Strength | Cost |
|---|---|---|
| Ignore duplicates | Ships today | Corrupt counts, double charges |
| Inbox | Clear for side effects | Extra table; must share the transaction with the effect |
| LSN or version upsert | Natural for projections | Needs per-key order |
| Kafka EOS transaction | Strong when the sink is Kafka | Complexity; small set of sinks |
| Redis SETNX dedupe | Quick | TTL races; another dual-write |
What Kafka EOS does not solve
An idempotent producer and a Kafka transaction coordinate records and offsets inside Kafka. They do not:
- Make an HTTP charge or an email send idempotent.
- Repair a key that changed, or a composite primary key you modeled wrong.
- Cover Elasticsearch, SQL, or object storage. Those sinks still need their own guard.
EOS read-process-write helps when the sink is another Kafka topic. Say that limit, then point at the messaging lesson. Do not recite producer epochs here.
Ordering is not global exactly-once
Per-key order plus an idempotent upsert is a correct projection of that entity. A globally ordered exactly-once stream across partitions is not what you get. Do not design a cross-partition invariant as if it were.
Multi-table snapshots can interleave. If the business rule spans tables, prefer one outbox domain event over stitching raw table topics in the consumer.
Ack order: write the sink (or the inbox) then advance the consumer offset. Ack first, crash second, and the event is gone. The inverse crash (apply, then crash before ack) redelivers, and the guard turns that into a no-op.
LSN guard and ack order (run this)
Press Run. Snippets must be self-contained — no network, files, or native modules.
Press Run. Snippets must be self-contained — no network, files, or native modules.
The second call is the same offset, already acked, so it does not apply. The third call is a new offset with the same event id. The inbox keeps the balance at 5. That is the crash-between-apply-and-ack case once the first apply committed and the offset did not.
Interview Q&A
Is Debezium exactly-once?
Answer
Capture through Connect is effectively at-least-once because offsets flush in batches. Effects are exactly-once only when the sink dedupes. Say both sentences.
Why do message keys matter for CDC?
Answer
They pin an entity to a partition, they make compaction correct, and they give you per-entity order. The hash details live on the partition-keys lesson. Here you choose the key: primary key or aggregate id.
What is a tombstone?
Answer
A record with the entity key and a null value. Compacted topics use it to drop prior values for that key. Sinks use the delete event to remove the document. One without the other leaves ghosts or resurrections.
Inbox versus version upsert?
Answer
Inbox when the effect is not associative: charge, email, shipment. Version or LSN upsert when the effect is a projection of the row and older events must not win. Both want a stable id minted before publish.
Can I ack before writing the sink?
Answer
No. That is loss. Commit the sink, or the inbox row together with the effect, then ack. A crash in between redelivers and must no-op.
How should snapshot duplicates be handled?
Answer
The same as any redelivery. op=r must be safe to apply twice, and a later streaming event with a higher LSN must be allowed to win.
When do Kafka transactions help a CDC pipeline?
Answer
When the sink is Kafka and you can wrap the read-process-write in a transaction. They are not a blanket solution for SQL, search, or HTTP. Limits are on the delivery-semantics lesson.
How do you test it?
Answer
Crash injection between apply and ack, plus a fixture that replays the same event id and a fixture with an older LSN. Assert one side effect and a projection that does not move backward.
What if two instances of the consumer both see the event?
Answer
An in-memory set will double-apply. The inbox or the LSN predicate has to be in a store they share, in the same transaction as the effect.
Where does the outbox event id come from?
Answer
From the outbox insert inside the business transaction, before the relay publishes. The consumer must not allocate a fresh id. That contract is Transactional Outbox & Inbox Patterns.
Pitfalls
Draw connector offset, consumer offset, projection LSN, and inbox. Kill the process after the projection write and before the ack. Write the second delivery’s LSN and the expected projection. Then change the effect to “send email” and show the inbox row that makes the second delivery skip.