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:
- In the same transaction as the business change, insert an event row into an
outboxtable:(id, aggregate_id, type, payload, created_at). - CDC (or a poller) reads the outbox table and publishes each row to the broker.
- 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.