SysDesignPrep.com
Study guide 93 of 183

Change data capture and the outbox pattern

Keeping caches, search indexes, warehouses and other services in sync with a database: dual writes and why they fail, log-based change data capture with Debezium, the transactional outbox, ordering, backfills, schema changes and delete handling.

Reading is half of it. See this used in a real interview: walk through Design Yelp →

Most systems keep the same data in several places: the primary database, a cache, a search index, an analytics warehouse, other microservices. Keeping them in sync is surprisingly hard. Change data capture (CDC) reads every committed change from the database's own log and streams it to whoever needs it. Together with the outbox pattern, it is the standard answer to "how does the search index stay up to date?" or "how do you publish an event reliably when the order is saved?".

The dual-write problem

The naive approach writes to the database, then writes to the other system:

db.save(order)
search.index(order)      // or kafka.publish(OrderPlaced)

Things go wrong:

  • The process crashes between the two writes: the database has the order, the index never does.
  • The second write fails and is retried later, after a newer update, so the index ends up with older data.
  • Two concurrent requests interleave their writes in different orders in the two systems.

No amount of retrying fixes this, because the two writes are not atomic. You need one source of truth and a reliable way to derive everything else from it.

Log-based CDC

Every database already writes changes to a log for replication and recovery: the Postgres write-ahead log, the MySQL binlog, MongoDB's oplog, DynamoDB Streams. A CDC connector (Debezium is the common open-source choice) reads that log and publishes each insert, update and delete to a stream such as Kafka, usually one topic per table, keyed by primary key.

Advantages:

  • Captures every committed change, including those made by scripts and other apps, and nothing that rolled back.
  • Order per row is preserved when events are keyed by primary key.
  • No changes to application code and little extra load on the database.

Consumers then update search indexes, caches, warehouses and other services, each at its own pace. See message queues and streams.

The transactional outbox

Raw CDC exposes your table structure as an event format, which couples consumers to your schema. The outbox pattern fixes this:

  1. In the same transaction as the business change, insert an event row into an outbox table: (id, aggregate_id, type, payload, created_at).
  2. CDC (or a poller) reads the outbox table and publishes each row to the broker.
  3. Consumers receive a clean, intentional domain event like OrderPlaced.

Because the event and the change commit atomically, there is never an order without its event or an event without its order. Delivery is at least once, so consumers must be idempotent, using the event id. See distributed transactions and idempotency.

Common uses

  • Search indexing: listings or restaurants changed in Postgres, indexed into Elasticsearch within seconds. See search and indexing and Design Yelp.
  • Cache invalidation: delete or refresh cache keys when rows change, instead of from application code.
  • Analytics: stream changes into the warehouse or data lake instead of nightly dumps. See OLTP versus OLAP.
  • Microservice integration: other services build their own read models from events. See event sourcing and CQRS.
  • Migrations: replicate from an old database to a new one while both run.

Practical concerns

  • Initial snapshot (backfill): a new consumer needs existing data, not just new changes. Connectors take a consistent snapshot, then switch to the log at the right position.
  • Log retention: if the connector stops longer than the database keeps its log, you must re-snapshot. Monitor connector lag.
  • Ordering: ordered per key, not globally. Consumers that join several tables must handle events arriving in different orders.
  • Deletes: emitted as delete events or tombstones; consumers must apply them, or deleted data lives on in the index.
  • Schema changes: a new column flows through automatically; a renamed or removed one can break consumers. Use a schema registry and compatible changes. See schema evolution and serialization.
  • Large transactions: a bulk update of millions of rows becomes millions of events; throttle or route bulk jobs separately.

Polling as a simpler alternative

For modest volumes, a job that polls WHERE updated_at > last_seen works. It misses hard deletes (use soft deletes), can miss rows committed out of timestamp order (overlap the window and deduplicate), and adds query load. Log-based CDC is more reliable when it matters.

In the interview

When a design has a secondary store (search, cache, warehouse, another service), say: "the database is the source of truth; changes flow out via CDC (or an outbox) into Kafka, and an indexer consumes them idempotently." It answers consistency, reliability and decoupling in one sentence.

Checklist

  • No dual writes from application code for derived data.
  • CDC from the database log, keyed by primary key.
  • Outbox table for domain events, written in the business transaction.
  • Idempotent consumers, delete handling, and connector lag monitoring.
  • Snapshot plus log for backfills; compatible schema changes.

Open in your browser to sign in

Google does not allow sign-in inside this app's built-in browser. Open this page in Safari and sign in there. The link opens this same page.

Tap the ⋯ or share button at the top or bottom of the screen, then Open in browser. Or copy the link and paste it into Safari.