Change data capture turns a database's own write-ahead log into a stream of events: every insert, update and delete, in commit order, available to other systems seconds after it happens. The idea is simple and the payoff is large: search indexes, caches, warehouses and downstream services can follow the source of truth without dual writes and without nightly batch exports.
Getting the first event out is the easy part. The hard part is running the pipeline for years, through restarts, failovers, schema changes, backfills and quiet weekends when the source barely changes. This article treats CDC as a streaming system and follows a log-based pipeline from PostgreSQL through Debezium and Kafka to the consumers, concentrating on positions, delivery guarantees, snapshots and operations. The comparison of capture methods and the general snapshot handoff is in the CDC architecture overview, and the PostgreSQL decoding machinery in the logical decoding article.
Why CDC is a streaming problem
A CDC pipeline has every property of a stream processor. It has a source position that must be persisted and resumed from. It has a delivery guarantee that depends on the order in which it commits output and position. It has backpressure: if the consumer side stalls, the source must hold data, and in log-based CDC the source holding data means the database keeping WAL on disk. And it has state that can drift, because a downstream copy built from events is only correct if no event was skipped or applied twice in a way that matters.
Thinking of it this way clarifies the design questions. Where does the position live? What happens on restart between writing an event and recording the position? How does a consumer make duplicates harmless? How is the stream bootstrapped from a table that already has a billion rows? Each section below answers one of these.
The source side: slot, publication and connector
In PostgreSQL the database end consists of a publication, which names the tables to capture, and a logical replication slot, which records how far the consumer has confirmed and therefore how much WAL the server must keep. The connector, here Debezium's PostgreSQL connector running inside a Kafka Connect cluster, opens a replication connection on that slot, receives decoded changes through the built-in pgoutput plugin, and turns each into a Kafka record.
{
"name": "orders-cdc",
"config": {
"connector.class": "io.debezium.connector.postgresql.PostgresConnector",
"plugin.name": "pgoutput",
"database.hostname": "orders-db.internal",
"database.port": "5432",
"database.user": "cdc_reader",
"database.password": "${file:/secrets/cdc.properties:password}",
"database.dbname": "orders",
"slot.name": "orders_cdc",
"publication.name": "orders_pub",
"topic.prefix": "orders",
"table.include.list": "public.orders,public.order_items",
"snapshot.mode": "initial",
"heartbeat.interval.ms": "10000",
"signal.data.collection": "public.debezium_signal",
"provide.transaction.metadata": "true"
}
}With this configuration changes to public.orders land in the topic orders.public.orders, keyed by the table's primary key. Give the connector a dedicated database role with replication rights and read access only to the captured tables, and put the slot and publication names under version control: an orphaned slot left behind by a deleted connector is one of the most common ways CDC takes a database down.
Positions and delivery guarantees
Two positions exist and they are not the same. The slot's confirmed flush LSN, kept by PostgreSQL, is the point before which the server may discard WAL. The connector's offset, kept by Kafka Connect in its offsets topic, is the point the connector will resume from after a restart. Debezium writes events to Kafka, Kafka Connect periodically commits the offset of the last event it has delivered, and the connector then confirms that LSN to the slot.
The order is what produces the guarantee. Events are produced before their offset is committed, so a crash between the two replays some events on restart: the classic at-least-once behaviour. Kafka Connect supports exactly-once delivery for source connectors when the worker is configured for it and the connector supports it, which writes records and offsets in one Kafka transaction; whether your Debezium connector and version support it is something to check in its documentation rather than assume. Even with it, exactly-once ends at the Kafka topic. A consumer that writes to a database or a search index must still be idempotent, so design for duplicates and treat exactly-once source delivery as a reduction in noise, not a correctness mechanism. The general argument is laid out in the exactly-once semantics article.
The change event is a changelog entry
Each Debezium value carries an envelope: before and after row images, an op code (c create, u update, d delete, r read during snapshot, t truncate), a source block with the LSN, transaction id, table and commit timestamp, and ts_ms for when the connector processed it. The record key is the primary key.
For updates, before is only complete when the table's replica identity is FULL; with the default identity PostgreSQL logs only the old key. Large values stored out of line (TOAST) that did not change are not sent at all under the default identity, and Debezium substitutes a placeholder value, __debezium_unavailable_value by default. A consumer that writes after blindly will overwrite a real document body with that placeholder. Either set REPLICA IDENTITY FULL on tables with large columns, at the cost of more WAL, or have the sink merge rather than replace those columns.
A delete produces two records by default: the delete event, then a tombstone with the same key and a null value. The tombstone is what lets Kafka log compaction eventually remove the key entirely, so a compacted CDC topic converges to one record per live row: a table you can replay. Compaction mechanics are covered in the log compaction guide.
Topics, keys and ordering
Keying by primary key means every change to a given row goes to one partition, so per-row order is preserved. Nothing else is. Two rows in the same table may be in different partitions, and an orders row and its order_items rows are in different topics entirely. A consumer that needs a transaction's full effect, for example to publish an order only once all its items exist, cannot get it by reading topics independently.
Two remedies exist. With provide.transaction.metadata enabled, Debezium adds the transaction id and per-transaction event counts to each event and writes BEGIN and END markers to a separate transaction topic; a consumer can buffer events until it has seen the declared count for a transaction. The simpler remedy is to avoid needing it: capture an outbox table written in the same transaction instead of the raw tables, so each business event is one row with one key. That pattern is described in the transactional outbox article and is usually the better choice when the consumers are other services rather than analytic copies.
Snapshots and backfills
A new pipeline must first copy the existing rows. With snapshot.mode set to initial the connector takes a consistent snapshot, emits every row as an r event, then streams from the slot position taken at the snapshot's start, so nothing is lost at the boundary. On a large table this snapshot can take hours, and if the connector restarts mid-way it starts again.
Incremental snapshots fix both problems. Based on the watermark technique from Netflix's DBLog paper, they read the table in primary-key chunks while streaming continues, writing low and high watermark markers to the signal table around each chunk; a streamed change that arrives between the markers for a key already read in the chunk supersedes the chunk's copy. Progress is stored in the offsets, so a restart resumes at the last chunk. You trigger one by inserting a row into the signal table:
INSERT INTO public.debezium_signal (id, type, data)
VALUES ('backfill-2026-10-01', 'execute-snapshot',
'{"data-collections": ["public.orders"], "type": "incremental"}');The same mechanism handles adding a table to an existing pipeline and re-seeding a downstream system that lost data, which makes the signal table one of the most useful operational tools in the pipeline. Protect it: anyone who can write to it can force a full re-read of a large table.
Consuming the changelog safely
Most consumers materialise the stream into another store with upserts. Because events can arrive twice, and after a backfill may arrive older than what is already stored, the sink should only apply an event that is newer than the row it already has. The source LSN is a natural version for PostgreSQL sources:
INSERT INTO orders_copy (id, status, total, deleted, src_lsn)
VALUES (%(id)s, %(status)s, %(total)s, %(deleted)s, %(lsn)s)
ON CONFLICT (id) DO UPDATE
SET status = EXCLUDED.status,
total = EXCLUDED.total,
deleted = EXCLUDED.deleted,
src_lsn = EXCLUDED.src_lsn
WHERE orders_copy.src_lsn < EXCLUDED.src_lsn;def apply(event, cur):
v = event.value
if v is None: # tombstone: compaction marker, nothing to apply
return
lsn = v["source"]["lsn"]
if v["op"] == "d":
row = dict(v["before"], deleted=True)
elif v["op"] in ("c", "u", "r"):
row = dict(v["after"], deleted=False)
else: # "t": truncate, handle explicitly
raise ValueError("truncate needs an operator decision")
cur.execute(UPSERT_SQL, dict(row, lsn=lsn))Deletes become soft-delete rows with a version rather than physical deletes, so a delayed duplicate of an older insert cannot resurrect the row. Commit the consumer's Kafka offset only after the database transaction commits. Snapshot r events carry the LSN of the snapshot point, so the guard also orders correctly against streamed changes. In a stream processor such as Flink, the same envelope maps onto a changelog of inserts, update-before, update-after and deletes, and the processor's state plays the role of the version guard.
A worked example
An orders service writes about 300 changes per second. Three consumers follow orders.public.orders: a search indexer, a cache invalidator that deletes order:{id} keys, and a warehouse loader. On Saturday the search cluster is down for four hours for an upgrade. The indexer stops consuming, Kafka retains the events, and the database is unaffected, because the connector keeps reading and confirming the slot; Kafka, not PostgreSQL, absorbs the outage. That decoupling is the point of putting a log between source and sinks.
On Sunday the Kafka Connect cluster itself is down for two hours. Now the slot stops advancing and PostgreSQL retains WAL: at 300 changes per second with an average of 1 KB of WAL each, around 2 GB accumulates. That is fine. A forgotten slot from a connector deleted a month ago, at the same rate, is about 780 GB of retained WAL and very likely a full disk. Set max_slot_wal_keep_size so the server invalidates a runaway slot instead of filling its disk, and alert on slot lag long before that limit.
Operational hazards
- Quiet sources: if the captured tables rarely change but the database is busy elsewhere, the slot's position does not advance and WAL builds up. The heartbeat setting makes the connector commit offsets periodically, and
heartbeat.action.querycan write to a heartbeat table so there is always a captured change. - Failover: logical slots historically lived only on the primary and were lost on promotion, forcing a re-snapshot. PostgreSQL 17 can synchronise failover-enabled logical slots to physical standbys; on older versions, plan the re-snapshot into the failover runbook.
- Schema changes: adding a nullable column is backward compatible; renaming or changing a type is not. Route schemas through a registry with a compatibility rule so a breaking change fails at the producer rather than in every consumer.
- Truncate: decide in advance whether a sink mirrors a truncate or ignores it; a warehouse that silently loses history is worse than an alert.
- Snapshot load: chunked snapshots still read the whole table; schedule large ones off-peak and watch replica lag if you snapshot from a standby.
What to measure
Four numbers cover most incidents. Slot lag on the database: pg_wal_lsn_diff(pg_current_wal_lsn(), confirmed_flush_lsn) from pg_replication_slots, in bytes. Connector lag behind the source, which Debezium exposes as the MilliSecondsBehindSource streaming metric. Consumer group lag per sink in Kafka. And end-to-end freshness, measured by writing a timestamped canary row every minute and alerting when a sink has not seen it within its freshness objective. The first catches disk risk, the second and third locate the stall, and the fourth is the only one users actually feel.
Trade-offs
Log-based CDC gives low latency, captures deletes, and puts almost no query load on the source, but it couples you to the database's replication internals and to the physical schema: a column rename upstream is a contract change downstream. Capturing an outbox table restores a deliberate contract at the cost of writing the events in application code. Running Kafka between source and sinks adds a cluster to operate but isolates the database from sink outages, which a direct database-to-sink connector cannot do.
What to do next
- Inventory every logical slot on every database and match each to a live connector; drop orphans.
- Set
max_slot_wal_keep_sizeand alert on slot lag in bytes well below it. - Enable heartbeats, with a heartbeat action query if captured tables can go quiet.
- Decide replica identity per table, and set
FULLwhere sinks need old values or large columns. - Make every sink idempotent with a version guard on the source LSN and soft deletes.
- Create the signal table and rehearse an incremental snapshot on a non-critical table.
- Write the failover procedure for slots on your PostgreSQL version.
- Add a canary row and an end-to-end freshness alert per sink.