Change Data Capture - Complete Deep Dive
Prerequisites: Write-Ahead Log, Message Queues, Idempotency Used in: Keeping a search index in sync, feeding a warehouse, cache invalidation, and microservice data replication
What is Change Data Capture?
Change Data Capture turns every insert, update, and delete committed to a database into a stream of events that other systems can consume. The database stays the single source of truth. Everything derived from it - a search index, a cache, a warehouse, another serviceβs local copy - follows the stream instead of being written to separately.
Real-world analogy: A bank does not phone every department each time a transaction clears. It keeps a ledger, and anyone who needs to know reads the ledger from wherever they left off. The ledger already exists because the bank needs it for its own accounting. CDC is the realisation that your database keeps the same ledger, for the same reason, and you are allowed to read it.
The Dual-Write Problem
Here is the code everyone writes first, and it is wrong in more ways than it looks.
public void updatePrice(long productId, int priceCents) {
postgres.update(productId, priceCents); // 1
elasticsearch.index(toDocument(productId)); // 2
}
flowchart LR
APP["Product Service"]:::service
PG[("Postgres<br/>system of record")]:::data
ES[("Elasticsearch<br/>search index")]:::data
DRIFT["Silent drift<br/>row updated index stale"]:::async
APP -->|"1. UPDATE products"| PG
APP -->|"2. Index document"| ES
ES -->|"Step 2 fails or times out"| DRIFT
PG -->|"No transaction spans both"| DRIFT
classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
classDef edge fill:#1e3a5f,stroke:#60a5fa,color:#e2e8f0
classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
classDef async fill:#3b1f5e,stroke:#c084fc,color:#e2e8f0
classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0
| Color | Meaning |
|---|---|
| π£ Purple | Clients |
| π’ Green | Services |
| π‘ Yellow | Data stores |
| π΅ Blue | Edge / CDN |
Every way this partially fails:
- Step 2 throws. Elasticsearch is mid rolling-restart, or a circuit breaker is open. The row is correct, the index is stale, and nothing in your monitoring knows the difference.
- It times out, but it actually succeeded. You retry. Here that is harmless because indexing by document ID is idempotent. Make step 2 βpublish an eventβ instead and you have now published it twice.
- The process is killed between 1 and 2. No exception, no log line, no retry. That rowβs index entry is wrong permanently.
- Someone reorders the calls. Index first, then the database write fails, and search now returns a product that does not exist.
- Two concurrent updates to the same row. The database serialises them to v1 then v2. The two index calls race over the network and land v2 then v1. The index now holds a version the database never had, and no retry will ever repair it, because at the moment of the last write nothing was wrong.
- You wrap it in a transaction. A Postgres transaction cannot roll back an HTTP call to Elasticsearch. Two independent systems do not share a commit.
The last point is the actual problem, and the first five are symptoms. You cannot make two systems atomic without either two-phase commit, which nobody wants on a request path, or a pattern that makes one systemβs commit cause the otherβs write. CDC is that pattern. So is the outbox.
What makes dual writes worse than an outage is that the damage is quiet. A write that fails loudly gets paged. A dual write that half-succeeds leaves a search index that is quietly wrong for some fraction of rows, forever, growing a little every week, and you find out about it from a support ticket.
Three Ways to Capture Changes
Polling an updated_at Column
SELECT * FROM products
WHERE updated_at > :last_seen
ORDER BY updated_at
LIMIT 1000;
Dead simple, no privileges needed, runs anywhere. It genuinely works for plenty of cases, and dismissing it is a mistake. What it gets wrong:
Deletes are invisible. A deleted row is not returned by any query, so a consumer has no way to learn it is gone. The workaround is soft deletes - a deleted_at column and never actually deleting - which means reshaping your data model to serve your pipeline.
Intermediate states vanish. Poll every 10 seconds. A row that goes pending to paid to shipped inside that window produces exactly one event: shipped. If a consumer exists to send a receipt on paid, it never fires. For a search index that only cares about current state, collapsing is a feature. For an event consumer it is a correctness bug.
The load is constant and mostly wasted. Every poll is a range scan on the primary whether anything changed or not, and it needs an index on updated_at, which taxes every write.
Timestamps lie about commit order. If the application sets updated_at, two app servers with 200ms of clock skew can commit out of order and you will skip a row. If now() sets it inside a transaction, the value is the transaction start time - so a long transaction can commit a row stamped earlier than a watermark you have already passed. The careful version uses a monotonic sequence plus a deliberate overlap window and relies on the consumer being idempotent.
Latency floors at the poll interval. Poll faster to improve it, pay more load.
Triggers Writing to an Audit Table
CREATE TRIGGER products_audit
AFTER INSERT OR UPDATE OR DELETE ON products
FOR EACH ROW EXECUTE FUNCTION write_change_row();
This captures everything polling misses: deletes, every intermediate state, and both before and after values. It is also transactional - the business change and its audit row commit together or not at all - which is a real guarantee polling cannot offer.
The costs are operational rather than theoretical.
The trigger runs inside your transaction, on the hot path. Every write becomes two writes plus a function call, which extends the transaction and doubles write volume. Bulk operations suffer most: a 100,000-row update becomes 200,000 row writes. The audit table then becomes a hot, permanently growing table needing its own retention job, vacuum attention, and index.
And the logic lives in the database. A PL/pgSQL function is invisible to your applicationβs code review, missing from your deploy pipeline, handled badly by migration tools, and the last place anyone looks when behaviour is strange. The characteristic failure is partial and silent: a developer adds a column, nobody updates the trigger function, and the audit stream quietly stops carrying that field.
Triggers earn their place when you need capture inside a database you cannot reconfigure - no replication slot available, no superuser, logical decoding disabled by the managed provider.
Log-Based CDC
Every durable database already maintains a complete, ordered record of every change, because it needs one for crash recovery and replication. Postgres calls it the WAL, MySQL the binlog, MongoDB the oplog, SQL Server the transaction log. That record is a better change stream than anything you could build, and it is already being written.
A log-based connector does not query your tables at all. It registers as a replication client and reads the log, the same way a read replica does.
Why this is the one that wins:
- No application change. Not a line. Capture happens below your code, which also means it captures writes from migrations, from a
psqlsession, from the ORM, and from that cron job nobody owns. - Deletes are first-class. A delete is a log record like any other.
- Ordering is correct for free. The log is in commit order by construction. No clocks involved.
- Every intermediate state survives.
pending,paid,shipped- three records, in order. - The source database barely notices. Reading the log is a sequential read of a file the database already wrote. No table scans, no extra index, no trigger on the write path.
Two costs, stated plainly. You need privileged configuration - wal_level = logical in Postgres, binlog_format = ROW in MySQL - and you now operate a connector plus its offset store. And the one that causes real incidents: an unconsumed replication slot stops the database recycling WAL segments. Leave the connector down over a weekend and Postgres will retain every segment since it stopped, until the disk fills. Alert on replication slot lag with the same seriousness you alert on disk space.
How Log-Based CDC Actually Works
flowchart LR
APP["Product Service"]:::service
PG[("Postgres<br/>system of record")]:::data
WAL[("WAL<br/>replication slot")]:::data
CONN["Debezium connector<br/>replication client"]:::async
K["Kafka topic<br/>one per table"]:::async
IDX["Search indexer<br/>idempotent upsert"]:::service
ES[("Elasticsearch")]:::data
APP -->|"1. Commit the row"| PG
PG -->|"2. Durable log record"| WAL
WAL -->|"3. Stream from last LSN"| CONN
CONN -->|"4. Before and after image"| K
K -->|"5. Consume and transform"| IDX
IDX -->|"6. Bulk index by document id"| ES
classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
classDef edge fill:#1e3a5f,stroke:#60a5fa,color:#e2e8f0
classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
classDef async fill:#3b1f5e,stroke:#c084fc,color:#e2e8f0
classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0
The connector is a replication client. It opens a replication connection and asks the database for changes starting from a position. Postgres calls that position an LSN - a log sequence number, effectively a byte offset into the WAL stream. MySQL uses a binlog filename plus offset, or a GTID. Logical decoding translates raw WAL records, which are about pages and bytes, into logical row changes; pgoutput ships with Postgres and is the usual plugin.
The log position is the checkpoint. The connector periodically commits the last processed LSN to durable storage - a Kafka topic, in Debeziumβs case. On restart it resumes from there. Separately, it confirms a position back to the database, and that is what permits Postgres to release WAL segments. Two different bookkeeping jobs, and both matter: the first prevents data loss, the second prevents a full disk.
Events carry before and after images. An update event contains the row as it was and as it now is, so a consumer can tell what changed without keeping its own previous copy.
{
"op": "u",
"ts_ms": 1727740800123,
"source": { "db": "shop", "table": "products", "lsn": 24857120, "txId": 94812 },
"before": { "id": 42, "name": "Kettle", "price_cents": 2999, "stock": 7 },
"after": { "id": 42, "name": "Kettle", "price_cents": 2499, "stock": 7 }
}
op is c for create, u for update, d for delete, and r for a row read during the initial snapshot. A delete arrives with a populated before, a null after, and then a second message with a null value - a tombstone - so a log-compacted topic can actually drop the key rather than retain a deletion marker forever.
π‘ In Postgres, before-images for updates and deletes need REPLICA IDENTITY FULL on the table. The default logs only the primary key, so a consumer that reads before.price_cents will find nothing until someone changes that setting - and changing it increases WAL volume for every write.
Snapshot Then Stream
Your products table already holds 50M rows and the WAL only goes back a few hours. Streaming alone hands a consumer the changes since you switched it on, which is not a usable search index. You need the existing rows first, and you need the handover to have neither a gap nor an unbounded lock.
The sequence:
- Record the current log position.
- Read the whole table, emitting each row as a synthetic
revent. - Begin streaming from the position recorded in step 1.
Done naively, step 2 holds a transaction open for hours, blocking DDL and bloating the WAL the entire time. What real connectors do instead:
Consistent snapshot without a long lock. Open a REPEATABLE READ transaction, capture the LSN, and read within that snapshot. MVCC gives a consistent view without blocking writers. A brief lock may still be needed to read the schema safely.
Overlap rather than gap. Start streaming from a position at or before the snapshotβs, so any change that happened during the snapshot is replayed. You get duplicates instead of gaps - acceptable, because the consumer has to be idempotent anyway and a row keyed by primary key is a last-write-wins upsert.
Incremental snapshots. Debezium can chunk a table by primary-key ranges and interleave those chunks with the live stream, so a 50M-row bootstrap does not block streaming for hours and survives a connector restart mid-way. The same mechanism lets you re-snapshot one table later without stopping everything, which you will want the first time a consumer loses data.
The principle worth stating out loud in an interview: choose duplicates over gaps, every time. A duplicate is a consumer-side problem with a known fix. A gap is silent, permanent, and discovered six months later.
Delivery Semantics: At-Least-Once, and Why
CDC is at-least-once. Not sometimes - structurally.
The connector emits events and then commits its offset. Crash in between, and those events are re-emitted on restart. Commit the offset first instead and you get at-most-once, which means gaps. There is no third option without a transaction spanning the database, the connector, and the broker. Add the deliberate snapshot overlap, any re-snapshot of a table, and producer retries on an ambiguous timeout, and duplicates are a normal operating condition rather than an incident.
So consumers must be idempotent. CDC makes that unusually easy, because every event carries a primary key and a monotonic log position. Two patterns cover nearly everything:
- Upsert keyed on the primary key. Applying the same row twice changes nothing. This is what a search indexer and a cache invalidator do, and it is why those are the easiest CDC consumers anyone ever builds.
- Version guard for effects you cannot repeat. Store the last applied LSN per key and discard any event at or below it. Needed when the consumer increments a counter, sends an email, or moves money - replaying those is not harmless.
Why exactly-once end to end is harder than vendors imply. There are four hops: database to connector, connector to broker, broker to consumer, and the consumerβs actual side effect. Kafka transactions can make the middle hops effectively-once within Kafka by committing the processed output and the consumer offset in one transaction. The moment the effect lands outside Kafka - an Elasticsearch bulk call, an HTTP request, a row in another database - no transaction covers it, and you are back to deduplicating in the consumer. βExactly-onceβ on a datasheet almost always means βexactly-once inside our system, assuming your sink is idempotent.β That is a true statement. It is not the one you read.
CDC Versus the Outbox Pattern
Both solve the dual-write problem the same way - the database commit becomes the single source of truth for what happened. They differ entirely in what crosses the boundary.
The outbox pattern gives you events you designed. Your service writes an OrderPlaced row into an outbox table in the same transaction as the order itself, with the fields you chose, the name you chose, and a version you control. Internal columns stay internal. Rename a column tomorrow and no consumer notices. The cost is application code: every event must be explicitly written, every new event is a code change, and a developer who forgets the outbox insert ships a silent gap.
CDC gives you every change for free. No application code, nothing to forget, complete by construction - including writes from scripts and migrations that no event-publishing code path would ever have covered. The cost is that your consumers are now coupled to your table schema. products.price_cents is in three other teamsβ pipelines, you did not know, and you can no longer rename it.
When each is right:
- CDC when the consumer wants current state rather than domain events: a search index, a cache invalidator, a replica in another store, a warehouse load, an analytics stream. The consumer wants the row, so the row is the correct contract.
- The outbox when the consumer wants domain events with business meaning -
PaymentCaptured,SubscriptionCancelled,FraudCheckPassed. These do not map one-to-one onto row changes: a single state transition may touch three tables, and the event name carries intent that no diff can express. Also use it whenever the consumer is another team, because an event is a contract you can version and a table is not. - Both, which is the common production answer. Run the outbox table itself through log-based CDC. You get designed, versioned events plus zero polling load and true commit ordering. Debezium ships an outbox event router for precisely this, and it is what the Digital Wallet and Notification System designs do.
In one line: CDC exports your data model. The outbox exports your domain.
Schema Evolution Will Break Your Consumers
ALTER TABLE products RENAME COLUMN price_cents TO price_minor_units;
A one-line migration. It passes review, deploys cleanly, and your service is entirely fine. Downstream:
- The search indexer reads
after.price_cents, gets null, and indexes every product at a price of zero. - Over in the warehouse, the load fails on an unknown column. Or worse, it succeeds with a null column, and the finance dashboard quietly reports nothing.
- Your cache invalidator is unaffected, because it only ever reads the key. Luck, not design.
Nobody was warned, because the contract was never written down anywhere. This is the real cost of CDC and it is not a configuration mistake.
Mitigations, roughly in order of how much they help:
A schema registry. Confluent Schema Registry, AWS Glue Schema Registry, or Apicurio holds a versioned schema per topic and rejects a producer whose schema violates the configured compatibility rule. Set BACKWARD compatibility and the registry refuses the rename at publish time instead of letting nine consumers discover it independently. The value is not that it prevents change - it converts a silent data-corruption bug into a loud pipeline failure.
Additive-only discipline at the source. Add the new column, backfill it, migrate consumers, drop the old one in a later release. Expand-contract migration, applied to a contract you did not realise you had.
A transform layer in front of consumers. Map raw CDC rows into a stable published schema inside a stream processor, so a rename blasts one job rather than nine teams. Note what this is: the outbox pattern, rebuilt downstream, after the fact.
Knowing who consumes you. Write down the topics and their consumers. Deeply unglamorous, and it prevents more incidents than the registry.
What no tool catches: a semantic change. Switch
price_centsfrom cents to paise and the schema stays valid, the registry approves it, every consumer keeps running, and every price downstream is wrong. Only a documented contract and a conversation prevent that one.
Comparison Table
| Β | Polling updated_at |
Triggers to audit table | Log-based CDC | Outbox |
|---|---|---|---|---|
| Captures deletes | No, needs soft deletes | Yes | Yes | Only what you emit |
| Intermediate states | Collapsed within the interval | Yes | Yes | Only what you emit |
| Ordering | Timestamp-based, skew-prone | Commit order per table | True commit order from the log | Commit order, via CDC or poll |
| Latency | The poll interval, 1s to minutes | Poll interval on the audit table | Sub-second | Sub-second with CDC |
| Source overhead | Repeated range scans on the primary | Extra write inside every transaction | Near zero, sequential log read | One extra insert per event |
| Application change | None | None, but database-side code | None | Every event is code |
| Consumer coupling | Your table schema | Your table schema | Your table schema | An event schema you own |
| Privileges needed | None | Create trigger | Replication slot plus config | None |
Common Interview Questions
Q: βWhy not just write to the database and Elasticsearch in the same request?β A: Because there is no transaction spanning them, so every partial failure leaves silent divergence - and the nastiest case is not a failure at all. Two concurrent updates to the same row serialise correctly in the database and then race over the network, so the index ends up holding a version the database never had, with nothing to retry. Make the database commit the trigger for the second write, through CDC or the outbox, and the ordering problem disappears because the log is already in commit order.
Q: βIs CDC exactly-once?β A: No. It is at-least-once, structurally: the connector emits then commits its offset, so a crash between the two re-emits. Flipping the order would give you at-most-once, which means gaps, and gaps are far worse. Build consumers that upsert on the primary key, or keep a last-applied LSN per key for effects you cannot repeat. Kafka transactions can get you effectively-once within Kafka, but the moment your sink is an HTTP call or another database, dedup is back on you.
Q: βHow do you bootstrap a table that already has 50M rows?β
A: Snapshot then stream. Record the log position, read the table inside a REPEATABLE READ transaction so MVCC gives a consistent view without blocking writers, emit each row as a read event, then stream from the recorded position. Overlap the handover deliberately so changes during the snapshot are replayed rather than lost. For a table that large, use incremental snapshots that chunk by primary-key range and interleave with the live stream, so the bootstrap is resumable and does not stall streaming for hours.
Q: βCDC or outbox for notifying another teamβs service?β
A: Outbox. Exposing your table schema to another team means they depend on price_cents and you can never rename it, and a row diff cannot carry intent - a status column flipping to cancelled loses the reason and the actor. Publish a versioned SubscriptionCancelled event you own. Then run the outbox table through log-based CDC so you still get sub-second delivery with commit ordering and no polling load.
Q: βWhat breaks if the CDC connector is down for two days?β A: Two things, in order of severity. The replication slot has not advanced, so Postgres has retained every WAL segment since the connector stopped - that is a disk-full incident on your primary, which is much worse than a stale index. And consumers are two days behind. Recovery is usually uneventful because the connector resumes from its committed LSN and catches up, provided the WAL is still there. If someone dropped the slot to reclaim the disk, you now have a gap and the fix is a re-snapshot. Alert on slot lag in bytes, not just on connector liveness.
Q: βHow do you ship a column rename safely once CDC consumers exist?β A: As an expand-contract migration against a contract you now treat as public. Add the new column, backfill it, write both for a release, migrate each consumer, then drop the old column. A schema registry with backward compatibility turns a mistake here into a rejected publish instead of nulls arriving in nine pipelines. And none of it protects you from a semantic change like switching currency units, which stays schema-valid and corrupts everything downstream.
When NOT to Use CDC
- You control both writers and a synchronous write is simpler. One service owning the Postgres row and the Redis key, in one request, with a retry and a TTL as the backstop - that is fine, and a CDC pipeline for it is over-engineering. The dual-write problem bites when the second write must not be lost. A cache entry with a 60-second TTL is allowed to be lost.
- Exposing the internal schema is unacceptable. If the consumer is another team, a partner, or anything with a contract, publish designed events through the outbox. Shipping your column names externally is shipping an API you cannot version and did not agree to support.
- The consumer needs domain events, not row diffs.
FraudCheckPassedis not derivable from a boolean flipping. The diff loses the reason, the actor, and the intent, and reconstructing them in the consumer puts your business logic in someone elseβs service. - You cannot operate the replication slot. A managed database with logical decoding disabled, no appetite for the WAL-retention failure mode, or nobody on call who will act on slot lag. Polling with soft deletes is the honest choice, and saying so is better than specifying infrastructure nobody will run.
- The derived store needs read-your-writes. CDC is asynchronous by definition. If the next read must observe the change, read the source of truth and let the derived store catch up.
- Your source is not a database at all. No log, nothing to tail. Capture at the application boundary instead.
| β Back to Fundamentals | Next: Outbox Pattern β |
Related Concepts
- Outbox Pattern β β the application-level alternative, and the thing you run CDC on top of
- Write-Ahead Log β β the log CDC reads, and why it already exists
- Idempotency β β mandatory for consumers, because CDC is at-least-once
- Event Sourcing β β where events are the source of truth rather than derived from rows
- Search Indexing β β the most common CDC consumer, and why upserts make it the easiest one
- Message Queues β β the transport, partitioning, and ordering guarantees underneath
Where This Shows Up
This concept is load-bearing in these designs - each link goes straight to the design that leans on it:
- Design an E-commerce Platform - catalogue writes fanning out to search, cache, and the warehouse from one committed stream
- Design Zomato - Debezium tailing the catalogue database into Elasticsearch, under 3 seconds end to end
- Design a Digital Wallet - Debezium on the Postgres WAL draining the outbox to Kafka beside a strict ledger
- Design a Notification System - outbox plus log-based capture so no notification is ever silently dropped
- Design a Q&A Forum - CDC keeping the search index in step with Postgres instead of dual-writing
- Design a Metrics Monitoring System - the same shape downstream: one stream feeding several derived stores independently
| All Concepts | All HLD Designs |