Debezium is an open-source set of change data capture connectors. Each connector reads a database's own replication stream (PostgreSQL's write-ahead log through logical decoding, MySQL's binary log, and equivalents for other databases) and turns every committed insert, update and delete into an event on a Kafka topic. Because it reads the log rather than polling tables, it captures deletes, sees every intermediate change, and adds little load to the database. Most deployments run Debezium as source connectors inside Kafka Connect, which supplies the worker processes, offset storage, configuration REST API and converters.
Why CDC belongs in a streaming architecture at all, and the general hazards of replication slots, positions and changelog consumption, are covered in the CDC architecture article. This page is the Debezium operator's manual: how to deploy and configure the connectors, what is inside each event, how snapshots and signals work, which single message transforms (SMTs) to use, and how to turn on exactly-once delivery. Property names and values were checked against the Debezium stable reference on 2026-10-02; check the reference for your exact version before copying configuration.
How the pieces fit
A Kafka Connect cluster is a group of worker JVMs that share three internal Kafka topics: one for connector configuration, one for status and one for source offsets. You register a connector by posting JSON to the REST API; Connect schedules its task on a worker. The PostgreSQL and MySQL connectors each run a single task, because a database log is a single ordered stream and splitting it would break ordering.
The task does two things in sequence. First, depending on the snapshot mode, it takes a consistent snapshot of the existing rows and emits each as a read event. Then it streams changes from the log position that corresponds to the snapshot, so no change is missed between the two. Periodically, Connect commits the task's source position (a PostgreSQL LSN, a MySQL binlog file and position or GTID set) to the offsets topic. On restart the connector reads that offset and resumes there. Anything emitted after the last committed offset and before a crash will be emitted again, which is why consumers must be idempotent unless you enable exactly-once support.
Preparing the databases
For PostgreSQL, set wal_level = logical, create a user with replication permission and SELECT on the captured tables, and use the built-in pgoutput plugin (plugin.name also accepts decoderbufs, which needs an extension). Debezium creates a replication slot named by slot.name and uses a publication named by publication.name. In production create the publication yourself, listing only the tables you capture, and let a database owner manage it. If you need the full previous row in update and delete events, set REPLICA IDENTITY FULL on those tables; with the default identity, the before image of a delete carries only the primary key and updates usually carry no before image at all.
For MySQL, enable the binary log in row format with full row images (binlog_format=ROW, binlog_row_image=FULL), enable GTIDs so the connector can follow a failover to another server, and keep binlogs long enough to cover the longest connector outage you intend to survive. The connector user typically needs SELECT, RELOAD, SHOW DATABASES, REPLICATION SLAVE and REPLICATION CLIENT, plus LOCK TABLES if snapshots must lock. Each connector also needs a numeric database.server.id that is unique among everything replicating from that server, because to MySQL it is a replica.
Registering a connector
A minimal PostgreSQL connector for an orders service, posted to the Connect REST API:
curl -X PUT http://connect:8083/connectors/orders-pg/config \
-H 'Content-Type: application/json' -d '{
"connector.class": "io.debezium.connector.postgresql.PostgresConnector",
"database.hostname": "orders-db.internal",
"database.port": "5432",
"database.user": "debezium",
"database.password": "${file:/secrets/pg.properties:password}",
"database.dbname": "orders",
"topic.prefix": "orders",
"plugin.name": "pgoutput",
"slot.name": "debezium_orders",
"publication.name": "dbz_orders",
"table.include.list": "public.orders,public.order_lines,public.outbox,public.debezium_signal",
"snapshot.mode": "initial",
"heartbeat.interval.ms": "10000",
"signal.data.collection": "public.debezium_signal",
"key.converter": "io.confluent.connect.avro.AvroConverter",
"value.converter": "io.confluent.connect.avro.AvroConverter",
"key.converter.schema.registry.url": "http://registry:8081",
"value.converter.schema.registry.url": "http://registry:8081"
}'Using PUT .../config rather than POST /connectors makes registration idempotent, which suits GitOps-style deployment. The password comes from a config provider rather than sitting in the JSON. Avro with a schema registry is one choice of converter; JSON with embedded schemas works but makes every message several times larger. See schema registry for how compatibility rules interact with table changes. A MySQL connector looks the same apart from the class, database.server.id, and two extra properties described below.
Topics, keys and the schema history topic
By default each captured table gets its own topic named <topic.prefix>.<schema>.<table>, for example orders.public.orders (MySQL uses the database name in place of the schema). The record key is the table's primary key, so all changes to one row land in one partition and stay in order relative to each other. Ordering across rows or tables is not preserved, which is the main thing consumers get wrong; Kafka ordering covers why. A table without a primary key produces records with a null key unless you configure a message key column.
MySQL's binlog contains row images without column names, so the connector must know the table structure as of each position in the log. It records every DDL statement it sees in a schema history topic, configured with schema.history.internal.kafka.topic and schema.history.internal.kafka.bootstrap.servers. On restart it replays that topic to rebuild the schema at its offset. The topic must have exactly one partition so its order is global, and infinite retention; if it is compacted or expires, the connector cannot restart safely. PostgreSQL does not need this topic because pgoutput sends relation metadata in the stream.
The change-event envelope
Every Debezium data event value has the same envelope. A trimmed update to an order, as JSON:
{
"before": { "id": 1042, "status": "PENDING", "total_cents": 4599 },
"after": { "id": 1042, "status": "PAID", "total_cents": 4599 },
"source": {
"connector": "postgresql", "name": "orders", "db": "orders",
"schema": "public", "table": "orders",
"txId": 88213, "lsn": 3527410992, "ts_ms": 1790913601000,
"snapshot": "false"
},
"op": "u",
"ts_ms": 1790913601412
}The op field says what happened: c create, u update, d delete, r read (emitted during snapshots), t truncate and m for logical decoding messages. For a create, before is null; for a delete, after is null, and by default the delete event is followed by a tombstone, a record with the same key and a null value, so that log compaction can eventually drop the key entirely (see log compaction). The source block identifies the origin and position. Its ts_ms is when the change happened in the database; the top-level ts_ms is when the connector processed it, so the difference between them is your capture lag for that event.
Snapshots: initial, none, and incremental
The snapshot.mode property decides what happens to existing data. initial (the default) snapshots once when no offset exists, then streams. no_data captures structure but no rows and starts streaming immediately, for when the history is already loaded elsewhere. initial_only snapshots and stops. when_needed snapshots if no offset exists or the stored offset is no longer available in the log. always snapshots on every start. There are also configuration_based and custom modes for finer control.
An initial snapshot of a large table is a long transaction that the connector cannot resume halfway, and new tables added later are not covered by it. Incremental snapshots solve both. The connector reads the table in primary-key chunks while streaming continues, de-duplicating chunk rows against concurrent changes so the stream stays correct, and it records progress in its offsets so a restart resumes at the last chunk. You trigger one by inserting a row into the signaling table, which the connector captures like any other change:
CREATE TABLE public.debezium_signal (
id VARCHAR(42) PRIMARY KEY,
type VARCHAR(32) NOT NULL,
data VARCHAR(2048) NULL
);
-- Backfill a newly captured table without stopping the connector.
INSERT INTO public.debezium_signal (id, type, data) VALUES (
'backfill-order-lines-2026-10-02',
'execute-snapshot',
'{"data-collections": ["public.order_lines"]}'
);The signaling table must be in the connector's include list and named in signal.data.collection. Signals can also arrive on a Kafka topic, through JMX or from a file, selected with signal.enabled.channels; the database channel is enabled by default. First add the new table to table.include.list (and, for PostgreSQL, to the publication) before signalling it; incremental is the default snapshot type for this signal.
Heartbeats and why idle connectors break databases
A PostgreSQL slot holds WAL until the connector confirms it. If the captured tables are quiet but the database is busy with other tables, the connector receives nothing to emit, commits no newer offset, and the slot pins WAL until the disk fills. heartbeat.interval.ms makes the connector emit a heartbeat record periodically, which gives Connect a newer offset to commit. Where the captured tables can stay idle for long periods, also configure a heartbeat action query that writes to a small captured table, so a real change flows through and the confirmed position advances. On MySQL the equivalent risk is the opposite: if the connector is down longer than binlog retention, its offset points at a purged file and it must snapshot again.
Single message transforms: unwrap and outbox
Many consumers do not want the envelope. The io.debezium.transforms.ExtractNewRecordState SMT replaces each record's value with its after state. Its delete.tombstone.handling.mode option controls deletes: tombstone (the default) keeps tombstones, drop removes deletes and tombstones, rewrite keeps a delete as a row with __deleted set to true, and rewrite-with-tombstone and delete-to-tombstone combine the behaviours. add.fields copies metadata such as op or source.ts_ms into fields prefixed with a double underscore. Unwrapping suits upsert sinks; keep the full envelope for audit trails and anything that needs the before image.
Capturing internal tables couples consumers to your schema. The outbox pattern avoids that: in the same transaction as the business change, the service inserts a row into an outbox table containing an explicit event. The io.debezium.transforms.outbox.EventRouter SMT turns outbox inserts into clean events. By default it expects columns id, aggregatetype, aggregateid, type and payload; it routes by aggregatetype to a topic named by route.topic.replacement (default outbox.event.${routedByValue}) and keys by aggregateid. The outbox is treated as an append-only queue: the router expects inserts and ignores deletes, so the service can delete old rows on a schedule without emitting anything.
"transforms": "outbox",
"transforms.outbox.type": "io.debezium.transforms.outbox.EventRouter",
"transforms.outbox.route.topic.replacement": "events.${routedByValue}",
"transforms.outbox.table.field.event.key": "aggregateid"
Exactly-once delivery
By default the pipeline is at-least-once: after a crash, events since the last offset commit are produced again. Kafka Connect 3.3.0 added exactly-once support for source connectors, which writes records and offsets in one Kafka transaction. Debezium's MariaDB, MongoDB, MySQL, Oracle, PostgreSQL and SQL Server connectors support it. Turn it on with exactly.once.source.support=enabled on every worker of a distributed Connect cluster and exactly.once.support=required on the connector, leaving transaction.boundary at its default of poll. Consumers must read with isolation.level=read_committed to benefit.
Two caveats. The guarantee covers the hop from Connect into Kafka, not your consumer's side effects, so idempotent sinks are still the most robust design (exactly-once semantics explains the end-to-end picture). And the Debezium documentation itself notes open issues in the Kafka transaction protocol that may affect the guarantee, citing KAFKA-17734, KAFKA-17754 and KAFKA-17582; treat exactly-once as a strong reduction in duplicates, and keep consumers idempotent anyway.
Failure modes and operating guidance
| Failure | Symptom | Prevention or fix |
|---|---|---|
| Slot pins WAL | PostgreSQL disk fills while the connector is down or idle | Heartbeats, alert on slot lag in bytes, cap with max_slot_wal_keep_size and accept a resnapshot |
| Binlog purged | MySQL connector fails on restart with a missing position | Retention longer than worst outage; snapshot mode when_needed |
| Schema history lost | MySQL connector cannot parse rows after restart | One partition, infinite retention, never compacted, back it up |
| Incompatible schema change | Converter rejects records; task fails | Additive changes only; registry compatibility checks in CI |
| Large initial snapshot | Hours-long transaction, no resume | Use no_data plus incremental snapshots |
| Duplicates after restart | Same event twice downstream | Idempotent consumers keyed by source position; exactly-once source support |
| Failover loses slot | New primary has no slot; changes skipped or resnapshot needed | Plan failover explicitly; PostgreSQL 17 can synchronise logical slots to standbys |
Monitor the connector's task state through the REST status endpoint, the lag between source.ts_ms and processing time, replication slot or binlog lag on the database, and the Connect worker's JMX metrics. Restart failed tasks automatically only after alerting a human, because a task that fails on a poison record will fail again.
What to do next
- Enable logical WAL or row binlogs with full images and GTIDs on a staging database, and create a dedicated connector user with the minimum privileges.
- Create the publication and signaling table yourself, then register a connector with PUT on its config endpoint, using a secret provider for credentials.
- For MySQL, create the schema history topic by hand with one partition and infinite retention before the first start.
- Configure heartbeats, alert on slot or binlog lag, and test what happens when the connector is stopped for a day.
- Decide envelope versus unwrapped values per consumer, and move cross-team events to an outbox table routed with EventRouter.
- Enable exactly-once source support on a Connect cluster at 3.3.0 or later, consume with read_committed, and keep sinks idempotent regardless.