Topic 10.4
Dual Writes, the Outbox Pattern and CDC
In one line
Writing to the database and publishing to Kafka (or updating a cache or search index) in two separate steps loses or duplicates events when either step fails. The outbox pattern writes the event into an outbox table in the same transaction as the business change; a relay or CDC (Debezium reading the WAL) publishes it. Consumers must handle duplicates and ordering per key.
Think of it like this
A shop that records a sale in the ledger and separately phones the warehouse. If the call drops, the sale exists but nothing ships. The outbox is writing "tell the warehouse" on the same ledger page as the sale; a runner reads the ledger and makes sure every note is delivered.
Key ideas
- 01
Dual-write failures: DB commit succeeds and the publish fails, so the event is lost; the publish succeeds and the commit fails, so it's a phantom event; two concurrent writers publish in a different order than they committed.
- 02
Outbox:
outbox(id, aggregate_type, aggregate_id, event_type, payload, created_at). The business write and the outbox insert commit together. A relay publishes rows in order and marks or deletes them, or CDC reads the inserts from the WAL. - 03
CDC with Debezium: a logical replication slot (
pgoutputplugin) streams row changes with their LSN; the Debezium connector (Kafka Connect) publishes them. The outbox event router SMT routes outbox rows to topics keyed by aggregate ID. It starts with an initial snapshot, then streams. - 04
Guarantees: at-least-once, so consumers see duplicates after connector restarts. Ordering is per key (the aggregate ID as Kafka key → one partition). Schema changes need compatible evolution (schema registry).
- 05
Operations: the replication slot retains WAL while the connector is down; cap it with
max_slot_wal_keep_sizeand alert. Clean the outbox (delete published rows or partition by day and drop). CDC use cases: search sync, cache invalidation, warehouse loading, audit, event-driven services. See the Kafka course for the consumer side.
Code & diagrams
CREATE TABLE outbox (
id uuid PRIMARY KEY DEFAULT gen_random_uuid(), -- uuidv7() on PG 18+
aggregate_type text NOT NULL,
aggregate_id text NOT NULL,
event_type text NOT NULL,
payload jsonb NOT NULL,
created_at timestamptz NOT NULL DEFAULT now()
);
BEGIN;
INSERT INTO orders (id, customer_id, total, status) VALUES (9001, 42, 1499.00, 'placed');
INSERT INTO outbox (aggregate_type, aggregate_id, event_type, payload)
VALUES ('order', '9001', 'OrderPlaced',
'{"orderId": 9001, "customerId": 42, "total": "1499.00"}');
COMMIT; -- both or neither
-- for Debezium CDC
ALTER SYSTEM SET wal_level = logical; -- needs restart
CREATE PUBLICATION outbox_pub FOR TABLE outbox;{
"name": "orders-outbox",
"config": {
"connector.class": "io.debezium.connector.postgresql.PostgresConnector",
"database.hostname": "pg-primary",
"database.dbname": "shop",
"plugin.name": "pgoutput",
"publication.name": "outbox_pub",
"slot.name": "orders_outbox_slot",
"table.include.list": "public.outbox",
"topic.prefix": "shop",
"transforms": "outbox",
"transforms.outbox.type": "io.debezium.transforms.outbox.EventRouter",
"transforms.outbox.route.by.field": "aggregate_type",
"transforms.outbox.table.field.event.key": "aggregate_id"
}
}Interview problem
The problem
Keep Elasticsearch in sync with PostgreSQL
Products are edited in PostgreSQL; the search index in Elasticsearch drifts (missing updates, deleted products still searchable). The app currently writes to both. Design a reliable sync, including reindexing.
When it breaks
Debezium connector down for a day
What you see
The replication slot retains WAL; the primary's disk fills and writes stop.
Fix & prevent
Set max_slot_wal_keep_size, alert on retained WAL and connector status, and have a runbook to re-snapshot if the slot must be dropped.
Outbox table never cleaned
What you see
It grows to billions of rows, its indexes bloat, and the relay's polling queries slow down.
Fix & prevent
Delete published rows in batches, or partition the outbox by day and drop old partitions.
Explain it without notes
Why can't you just publish to Kafka after the database commit?
Practice
Write the consumer-side logic that makes OrderPlaced processing idempotent.
Trade-offs
- ↔
Outbox + CDC gives reliable, ordered-per-key event publishing at the cost of extra infrastructure, WAL retention risk, and eventual consistency.
Done when you can
I can explain why dual writes fail and implement outbox + CDC with idempotent consumers.