Change data capture (CDC) turns every committed insert, update and delete in a database into an event that other systems can consume. It exists to solve the dual-write problem. If a service writes a row to its database and then publishes a message, one of the two will eventually fail without the other, and a search index, cache or warehouse silently drifts from the truth. CDC removes the second write. The database commits once, and the change stream is derived from what actually committed.

This article covers CDC as an architecture that works across database engines: how to capture changes, what an event looks like, how to take a consistent initial copy and then switch to streaming, how to keep order, handle schema changes and deletes, make sinks safe to replay, and run the thing without filling the database's disk. For how PostgreSQL decodes its write-ahead log internally, see CDC with logical decoding.

Advertisement

Three ways to capture changes

There are three capture methods, and they differ in what they can see and what they cost the source.

MethodHow it worksSees deletes?Cost and limits
Query-basedPoll WHERE updated_at > :lastNo, unless soft-deletedMisses intermediate updates and rows with clock skew; adds read load; needs an indexed timestamp on every table
Trigger-basedTriggers copy each change into a changelog tableYesAdds write latency to every transaction; triggers must be maintained through schema changes
Log-basedRead the database's own transaction logYes, with before-imagesNear-zero write overhead; needs log access, retention management and per-engine setup

Log-based capture is the default for serious systems because the transaction log already contains every committed change, in commit order, including deletes, and reading it does not touch the tables. Each engine exposes it differently. PostgreSQL needs wal_level=logical, a replication slot and a publication. MySQL needs binlog_format=ROW and, if you want full before-images, binlog_row_image=FULL. Oracle and SQL Server have their own log-mining or change-table mechanisms. A connector framework such as Debezium hides these differences behind one event format.

The architecture end to end

Log-based CDC: the database log is the source, the connector turns it into keyed eventsSource databaseOLTP writesTransaction logWAL / binlog / redoCDC connectorsnapshot + streamSlot / positionwhat is confirmedEvent log (Kafka)topic per table, key = PKSchema registryversioned schemasSearch indexupsert by keyWarehousemerge by positionCache invalidationdelete by keycommitdecodeack LSNeventsregisterThe connector acknowledges a log position only after the broker has the events; until then the database keeps the log.Consumers keep their own offsets, so each sink can replay from the event log without touching the source again.
Log-based CDC pipeline. The event log decouples sinks from the source, and the source retains its log until the connector confirms a position.

The connector reads the log, decodes row changes, and writes one event per changed row to a topic per table, keyed by primary key. Only after the broker has durably accepted the events does the connector confirm the log position back to the database. That ordering is what makes the pipeline at-least-once. A crash between publish and confirm means the connector restarts from the last confirmed position and republishes some events. Nothing is lost, but duplicates are normal, and every sink must tolerate them.

The event log in the middle is what makes CDC an architecture instead of a point-to-point sync. Each sink keeps its own offset. A new sink can start from the beginning of a compacted topic, and a broken sink can rewind, without asking the source database for anything.

Advertisement

Anatomy of a change event

Debezium's envelope has become the de facto shape. This example assumes REPLICA IDENTITY FULL on the orders table.

{
  "before": {"id": 1182, "status": "PENDING", "total": 49.00},
  "after":  {"id": 1182, "status": "REFUNDED", "total": 49.00},
  "source": {"connector": "postgresql", "db": "shop", "schema": "public",
             "table": "orders", "lsn": 24023128, "txId": 5521, "ts_ms": 1790841600000},
  "op": "u",
  "ts_ms": 1790841600412
}

op is c for create, u for update, d for delete and r for a read during a snapshot. before and after are row images. before is null on create, and after is null on delete. source records where the change came from, including the log position (lsn here, a binlog file and offset on MySQL) and the transaction id. The position is the most useful field for sinks, because it is monotonic per source and lets a sink tell an old version from a new one.

Whether before contains the full old row depends on the source. In PostgreSQL it contains only the key columns unless the table's replica identity is set to FULL. Decide per table whether consumers need old values, for example to compute a diff or to route on a changed column, because full images cost log volume.

Snapshot, then stream: the handoff

A new pipeline has to copy existing rows before it can stream changes, and the copy and the stream must line up with no gap and no lost update. The classic approach records the current log position, reads every table inside one consistent snapshot transaction, emits those rows as r events, and then streams from the recorded position. Changes that happened during the copy appear in the stream after their snapshot rows, so the final state is correct.

That approach has two problems at scale. The snapshot transaction can run for hours on a large table, which holds back vacuum or undo cleanup and grows the log. And streaming cannot start until the copy finishes. Netflix's DBLog design solves both with watermarks. The connector streams continuously and copies the table in primary-key chunks. For each chunk it writes a low watermark to a signal table, selects the chunk, and writes a high watermark. When it sees the watermarks come through the log, it drops any chunk row that was also changed in the stream between them, since the stream already carries a newer version, and emits the rest. No long transaction is needed, streaming never pauses, and a backfill can restart at the last finished chunk.

Debezium implements this as incremental snapshots, triggered through a signal table that is itself captured:

-- One-time setup, in the source database
CREATE TABLE debezium_signal (id VARCHAR(42) PRIMARY KEY, type VARCHAR(32) NOT NULL, data VARCHAR(2048) NULL);

-- Later: backfill one table while streaming continues
INSERT INTO debezium_signal (id, type, data) VALUES (
  'backfill-orders-2026-10-01',
  'execute-snapshot',
  '{"data-collections": ["public.orders"], "type": "INCREMENTAL"}'
);

Signals can also arrive over Kafka, JMX or a file if the connector enables those channels in signal.enabled.channels. Use incremental snapshots to add a table to a running pipeline or to repair a sink, without re-running a blocking snapshot of everything.

Ordering, keys and transactions

The log is totally ordered, but the event log is not. Kafka orders messages only within a partition. Keying each event by primary key sends all changes to one row to one partition, so they stay in order per row. That is the guarantee most sinks need: the last event for a key is the current state. Changes to different rows in different partitions can be consumed in any interleaving, and changes to different tables are in different topics.

Two consequences are easy to miss. First, changing the partition count of a topic changes which partition a key maps to, and can briefly reorder a key's events. Plan partition counts up front. Second, a sink that needs transactional consistency across tables, such as an order and its line items appearing together, cannot get it from per-row ordering. Either join and buffer by transaction id in a stream processor, using the transaction metadata the connector can emit, or accept that readers of the sink will see brief partial states. Before building transactional reassembly, check whether the consumers really need it.

Primary keys matter more than usual. A table without a primary key cannot be keyed, so its events cannot be ordered or upserted reliably. Add a key or leave the table out.

Schema evolution

Every ALTER TABLE on the source changes the shape of its events. The connector reads the new schema from the log and, with a schema registry, registers a new schema version, so consumers can decode old and new events side by side. Configure compatibility so the registry rejects changes that would break readers. Adding a nullable column is backward compatible. Renaming or dropping a column, or narrowing a type, is not.

The operational rule is that the source schema is now a public API. Agree on expand-and-contract migrations: add the new column, backfill it, move consumers, then drop the old column in a later release. Sinks with strict schemas, like warehouse tables, need an automated path for added columns and an alert, not an automatic change, for everything else.

Deletes and tombstones

A delete produces a d event whose after is null. With tombstones.on.delete enabled, the connector follows it with a tombstone, a message with the same key and a null value. The tombstone exists for log compaction: a compacted topic keeps only the latest message per key, and a tombstone tells the compactor it may eventually remove the key altogether. Without tombstones, a compacted topic keeps the last pre-delete version forever, and a sink rebuilding from it resurrects deleted rows.

Sinks must handle deletes explicitly. An upsert sink that ignores null values will simply never delete. A sink that hard-deletes can be resurrected by a late, older update arriving after the delete during replay. The safer pattern is a soft delete that keeps the log position, shown in the next section, followed by periodic cleanup of rows deleted long ago.

Idempotent sinks: the worked example

Because delivery is at-least-once, the sink decides whether duplicates and reordering are harmless. The pattern that works for almost any keyed store is a conditional upsert on the source position. Take the order-refund event above: order 1182 moves from PENDING to REFUNDED at LSN 24023128. The sink applies it only if the stored version is older.

-- Target table carries the source position of the version it holds
CREATE TABLE orders_replica (
  id        BIGINT PRIMARY KEY,
  status    TEXT,
  total     NUMERIC(12,2),
  deleted   BOOLEAN NOT NULL DEFAULT false,
  src_lsn   BIGINT NOT NULL
);

-- Apply a create, update, snapshot read or delete. Replays and reordering are harmless:
-- an older position never overwrites a newer one, and a delete keeps its position.
INSERT INTO orders_replica AS t (id, status, total, deleted, src_lsn)
VALUES (:id, :status, :total, :is_delete, :lsn)
ON CONFLICT (id) DO UPDATE
   SET status = EXCLUDED.status, total = EXCLUDED.total,
       deleted = EXCLUDED.deleted, src_lsn = EXCLUDED.src_lsn
 WHERE t.src_lsn < EXCLUDED.src_lsn;

Replay the same event and the WHERE clause makes it a no-op. Deliver an older update after it, perhaps from a rewound consumer, and it is ignored too. Delete the order at a later LSN and the row becomes deleted = true with that position, so no older upsert can bring it back. Snapshot rows need care here: in an incremental snapshot, use the position at which the chunk was read so that a snapshot row never overwrites a newer streamed change. Positions are only comparable within one source, so a sink fed from several databases stores a position per source. Exactly-once transports, discussed in exactly-once streaming, reduce duplicates inside Kafka, but they do not remove the need for an idempotent write at the edge.

For CDC events you publish as a domain contract rather than raw table rows, publish through an outbox table and capture that instead. The outbox pattern lets you shape events deliberately rather than exposing your table layout.

Operating the pipeline

Here is a typical PostgreSQL connector configuration, with the properties that matter for operations:

{
  "connector.class": "io.debezium.connector.postgresql.PostgresConnector",
  "database.hostname": "orders-db.internal",
  "database.dbname": "shop",
  "plugin.name": "pgoutput",
  "slot.name": "cdc_shop",
  "publication.name": "cdc_shop_pub",
  "topic.prefix": "shop",
  "table.include.list": "public.orders,public.customers,public.debezium_signal",
  "snapshot.mode": "initial",
  "signal.data.collection": "public.debezium_signal",
  "heartbeat.interval.ms": "10000",
  "heartbeat.action.query": "UPDATE cdc_heartbeat SET ts = now() WHERE id = 1",
  "tombstones.on.delete": "true"
}

The biggest operational risk is log retention. A replication slot makes PostgreSQL keep WAL until the connector confirms it. If the connector stops, the WAL grows until the disk fills and the primary goes down. Set max_slot_wal_keep_size (PostgreSQL 13 and later) so that the database invalidates the slot instead of dying. A lost slot means a re-snapshot, which is much better than an outage. Alert on slot lag in bytes well before that limit. MySQL has the opposite failure: binlogs expire on their own schedule, and a connector that was down longer than the retention period cannot resume.

The quiet-table trap needs a heartbeat. If the captured tables are idle but the database is busy, the connector receives no events for its tables, never confirms a newer position, and WAL accumulates. heartbeat.interval.ms makes the connector commit its position periodically, and heartbeat.action.query writes a row the connector does see.

FailureSymptomResponse
Connector downSlot lag and disk use climbingPage on lag; cap retention with max_slot_wal_keep_size
Slot invalidated or binlog expiredConnector cannot resumeRecreate slot, incremental snapshot, rely on idempotent sinks
Incompatible schema changeRegistry rejects schema; connector stopsExpand-and-contract migration; fix and restart
Sink falls behindConsumer lag grows; source unaffectedScale consumers per partition; the log absorbs the backlog
Resurrected rowsDeleted records reappear after replayTombstones on; soft deletes with position guard

Also plan for failover. A replication slot lives on one server, and whether it survives promotion of a replica depends on your PostgreSQL version and configuration. Test a failover with the connector running before you rely on it.

What to do next

  1. List every consumer that currently dual-writes or polls, and pick the source tables that feed them.
  2. Enable log-based capture: wal_level=logical or row binlogs, and give every captured table a primary key.
  3. Decide replica identity or full row images per table, based on whether consumers need old values.
  4. Key topics by primary key, fix partition counts up front, and register schemas with a compatibility mode.
  5. Build every sink as a conditional upsert on source position, with soft deletes and tombstones enabled.
  6. Set max_slot_wal_keep_size, configure a heartbeat, and alert on slot lag and consumer lag.
  7. Rehearse recovery: invalidate a slot in staging and repair the sink with an incremental snapshot.
Key takeaway: CDC replaces fragile dual writes with a stream derived from what the database actually committed. Capture from the transaction log, key events by primary key so per-row order holds, and use watermark-based incremental snapshots to copy tables without long transactions or paused streaming. Treat the source schema as a public API, enable tombstones so compaction and deletes work, and make every sink a conditional upsert on the source position so duplicates, replays and reordering are harmless. Operationally, the log is both the source of truth and the main risk: cap slot retention, run heartbeats, watch lag, and rehearse recovery before you need it.