Command Palette

Search for a command to run...

Hectal
PHASE 10Advanced ~9 min· topic 4 of 5

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.

0/5 · 0%

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

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

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

  3. 03

    CDC with Debezium: a logical replication slot (pgoutput plugin) 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.

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

  5. 05

    Operations: the replication slot retains WAL while the connector is down; cap it with max_slot_wal_keep_size and 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

outbox.sqlsql
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;
debezium-outbox.jsonjson
{
  "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"
  }
}
outbox.mermaiddiagram
Rendering diagram…

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

01

Why can't you just publish to Kafka after the database commit?

Practice

01

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.