Build your own CDC pipeline (Debezium + outbox)
Your service writes to its DB and publishes to Kafka — and any crash between those two writes is permanent inconsistency. Build a Change Data Capture pipeline (modeled on Debezium + the outbox pattern) that closes the gap by making the database itself the event source.
- Scenes
- 12 interactive scenes
- Time
- about 84 minutes
- Topic
- Queues, Pub/Sub & Event Streaming
What you are building, and why
You have a service. It writes to its database and it publishes to a queue — and you have been bitten (or nearly bitten) by the gap between those two writes. The database commit succeeded but the queue publish failed; the publish ack got lost so a retry duplicated the event; two writers raced and the queue's order disagreed with the database's last-writer-wins. Each is a permanent inconsistency: a state the system cannot heal on its own, because the system is two systems with no shared atomicity boundary.
This curriculum is about closing that gap. The fix is Change Data Capture — turning the database's own write log into a stream of change events — and the engineering pattern that makes it land in production: the outbox pattern, where every business write also writes a row whose entire purpose is to be published. Modeled on Debezium and Postgres logical replication.
You start from the user's mental model. "I have a service that writes to its DB. I want other services to know about my changes without breaking my service." From that one frame, every concept arrives motivated: the WAL exists because the database needs it for recovery; the replication slot exists because the database needs to know how far the consumer has read; the snapshot exists because the WAL doesn't go back forever; the outbox exists because the service still wants atomicity. Eleven scenes plus a capstone, no slideshow.
The single insight that unlocks the design: the database's own write log is the source of truth your downstream consumers wanted all along — the outbox makes that source carry your domain events instead of your row deltas.
What you will be able to explain afterwards
- the dual-write problem
- polling CDC vs log-based CDC
- WAL / binlog as the event source
- Debezium connector + change events
- replication slot + LSN
- consistent snapshot, then stream
- schema registry + compatibility modes
- outbox pattern (atomic event publish)
- outbox cleanup (tombstone + log compaction, partition-drop, INSERT+DELETE same-tx)
- partition-by-aggregate ordering
- at-least-once + idempotent consumer
- 01The dual-write trap — no atomicity boundary across DB and KafkaService writes to its DB and publishes to Kafka — and any crash between those two writes is permanent inconsistency. Four scenarios, four divergences, one structural fix.~7 min
- 02Polling CDC — the lossy fixSELECT WHERE updated_at > last_seen is technically Change Data Capture but cannot see deletes, collapses intra-interval flips, and trades latency against DB load.~7 min
- 03The DB already has a log — reading the Postgres WAL and MySQL binlogPostgres WAL, MySQL binlog — the database already keeps a durable, ordered log of every change for replication. CDC reads this log instead of the tables.~7 min
- 04Debezium tails the logDebezium is a connector that registers as a replica, decodes each WAL/binlog record into a structured change event, and emits one Kafka message per row change.~7 min
- 05Replication slot — the bookmark that fills disksA replication slot is a server-side cursor identified by an LSN that stops Postgres from recycling WAL the connector hasn't read — and an inactive slot is the #1 Debezium production failure.~7 min
- 06Snapshot, then streamDebezium runs a consistent snapshot (op=r events), records the LSN at snapshot time, then switches to streaming from that LSN — so history and live updates stitch at one seam.~7 min
- 07Schema evolution leaks the tableRaw CDC inherits the source DDL — an ALTER TABLE RENAME COLUMN silently breaks downstream unless a Schema Registry enforces a compatibility mode that rejects the change at registration.~7 min
- 08The outbox is a contract, not a table — the transactional outbox and Debezium's Outbox Event RouterWrite the event into a dedicated outbox table inside the same transaction as the business write — the DB transaction makes both atomic, and CDC tailing the outbox emits domain events decoupled from the business tables.~7 min
- 09Outbox cleanup — pick your poisonHard-delete after emit, tombstone + log compaction, partition-by-date drop — three trade-offs across WAL noise, race risk, and operational complexity. INSERT+DELETE same-tx collapses cleanup into the write.~7 min
- 10Ordering is per key, never global — hash partitioning by primary key across Kafka partitionsKafka guarantees order within a partition; partition by aggregate_id keeps per-aggregate events ordered while accepting that cross-aggregate order is never preserved.~7 min
- 11At-least-once is the ceilingDebezium 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.~7 min
- 12Design canvas — compose your CDC pipelineCapstone: pick a workload (search index / audit log / microservice events / read model), configure the six slots, fire failures. The verifier traces each absorbed or broken outcome back to the scene that introduced the responsible component.~7 min
More in Queues, Pub/Sub & Event Streaming
Moving events between services without losing them: logs, work queues, fanout, change data capture, stream processing, and the delivery systems built on top.
- Build Build KafkaA partitioned, replicated, append-only log. The log is the database — internalize that, and a dozen product designs get easier.
- Build Build a Message Queue (RabbitMQ / SQS)A point-to-point work queue — the messaging primitive Kafka is NOT. Each message goes to one consumer, ack deletes, retries push to a dead-letter queue, and a poisoned message is everyone's problem. Internalize ack vs visibility timeout vs DLQ vs prefetch vs FIFO groups — and learn to tell when Kafka is the wrong tool and when a queue is.
- Ad Click AggregatorStream processing with watermarks and exactly-once.
- Notification SystemPush, email, SMS. Idempotent. Failover.
- Distributed Cron — Mass Scheduled EmailSingle trigger, 50M recipient idempotency, catch-up.
Prefer to design it yourself?
The same subject as a staged workspace: draw the architecture, and a simulator traces requests through the boxes you drew.
Open the Build a CDC pipeline (Debezium + outbox) workspace