Kafka courseLesson 4 of 10
Kafka course · Lesson 4 of 10
Kafka Connect and Debezium CDC
How Kafka Connect moves data in and out of Kafka: source and sink connectors, workers, converters, SMTs, offsets and dead-letter queues, plus Debezium change data capture from Postgres.
On this page
- Sample setup
- Connect framework
- The pieces
- Worked example: the REST API
- Pitfalls
- In interviews
- Source vs sink connectors
- Worked example: a sink connector is a consumer group
- Pitfalls
- In interviews
- Standalone vs distributed
- Worked example: standalone mode
- Pitfalls
- In interviews
- Converters and SMTs
- How they fit together
- Worked example: unwrapping change events into rows
- The schema trap
- Pitfalls
- In interviews
- Offset storage
- Worked example: reading offsets
- Resetting or editing offsets
- Pitfalls
- In interviews
- Dead letter queue
- Worked example
- Pitfalls
- In interviews
- Debezium CDC
- How the PostgreSQL connector works
- Worked example: capturing changes
- Worked example: applying change events to a replica
- Snapshot modes
- Pitfalls
- In interviews
- Practice questions
- Key takeaways
Most data that lands in Kafka comes from a database, and most data that leaves Kafka goes into a warehouse, lake, search index or another database. Writing a producer or consumer for each of those is a lot of repetitive, failure-prone code. Kafka Connect is Kafka’s framework for that integration work, and Debezium is the most widely used set of Connect source connectors for change data capture (CDC). This lesson runs both against a real cluster.
Sample setup
The examples use:
- a three-broker Kafka 4.3.1 cluster on
localhost:19092; - a Kafka Connect worker in distributed mode with its REST API on
localhost:8083; - a PostgreSQL 16 database
shop(port 5499) started withwal_level = logical, with one table:
CREATE TABLE customers (id INT PRIMARY KEY, email TEXT NOT NULL, tier TEXT NOT NULL);
INSERT INTO customers VALUES (1, '[email protected]', 'bronze'), (2, '[email protected]', 'silver');
SELECT * FROM customers ORDER BY id;
id | email | tier
----+------------------+--------
1 | [email protected] | bronze
2 | [email protected] | silver
The worker configuration (worker.properties):
bootstrap.servers=localhost:19092,localhost:29092,localhost:39092
group.id=connect-cluster
key.converter=org.apache.kafka.connect.json.JsonConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
key.converter.schemas.enable=false
value.converter.schemas.enable=false
config.storage.topic=connect-configs
offset.storage.topic=connect-offsets
status.storage.topic=connect-status
config.storage.replication.factor=3
offset.storage.replication.factor=3
status.storage.replication.factor=3
listeners=http://localhost:8083
plugin.path=/opt/kafka/libs/connect-file-4.3.1.jar,/opt/connect-plugins
The Debezium PostgreSQL connector plugin (a folder of JARs) was unpacked into /opt/connect-plugins, then the worker was started with bin/connect-distributed.sh worker.properties.
Connect framework
What it is. Kafka Connect is a runtime for connectors: reusable plugins that copy data between Kafka and another system. You do not write a poll loop; you submit a JSON configuration to a Connect cluster, and the framework runs it, scales it, tracks its offsets and restarts it on failure.
The pieces
| Piece | Role |
|---|---|
| Worker | a JVM process running the Connect runtime; several workers form a Connect cluster |
| Connector | the configured job, for example “CDC from shop”; it splits work into tasks |
| Task | the unit of parallelism that actually moves data (tasks.max caps how many) |
| Converter | turns Connect’s internal records into bytes and back (JSON, Avro, Protobuf, String) |
| Transform (SMT) | a small per-record change applied between connector and converter |
| Plugin path | where workers find connector, converter and transform JARs (plugin.path) |
Worked example: the REST API
curl -s localhost:8083/
curl -s localhost:8083/connector-plugins
{"version":"4.3.1","commit":"26b251a451ce941d","kafka_cluster_id":"dnXyAz12QlOsbBQ3-avgJg"}
"class": "org.apache.kafka.connect.file.FileStreamSinkConnector",
"class": "io.debezium.connector.postgresql.PostgresConnector",
"class": "org.apache.kafka.connect.file.FileStreamSourceConnector",
"class": "org.apache.kafka.connect.mirror.MirrorCheckpointConnector",
"class": "org.apache.kafka.connect.mirror.MirrorHeartbeatConnector",
"class": "org.apache.kafka.connect.mirror.MirrorSourceConnector",
(The plugin list is filtered to the class names.) The three Mirror connectors ship with Kafka and power MirrorMaker 2. Common lifecycle calls:
| Call | Effect |
|---|---|
POST /connectors |
create a connector from {"name": ..., "config": {...}} |
PUT /connectors/{name}/config |
create or update (idempotent; good for automation) |
GET /connectors/{name}/status |
connector and task states, with stack traces for failed tasks |
PUT /connectors/{name}/pause / resume |
stop and restart processing, keeping tasks allocated |
PUT /connectors/{name}/stop |
stop and release tasks (required before editing offsets) |
POST /connectors/{name}/restart?includeTasks=true&onlyFailed=true |
restart failed tasks |
GET, PATCH, DELETE /connectors/{name}/offsets |
read, alter or reset offsets |
Pitfalls
- A
RUNNINGconnector with aFAILEDtask. The connector state alone is not health. Monitor task states. - Plugins are code. Connect runs whatever is on the plugin path inside the worker JVM; install only trusted connectors, and isolate plugins in their own folders to avoid dependency clashes.
- The REST API is unauthenticated by default. Restrict it to the network you trust, or enable authentication.
In interviews
“Why use Connect instead of writing a consumer?” Reuse of tested connectors, offset tracking, scaling across workers, restarts, standard error handling and DLQs, and configuration instead of code. Mention when not to use it: complex business logic belongs in a stream processor.
Source vs sink connectors
What it is. A source connector reads an external system and writes to Kafka topics. A sink connector reads Kafka topics and writes to an external system.
| Source | Sink | |
|---|---|---|
| Direction | system to Kafka | Kafka to system |
| Examples | Debezium (CDC from databases), JDBC source, S3 or file sources, MirrorMaker 2 | JDBC sink, S3/GCS sink, Elasticsearch/OpenSearch, Snowflake, BigQuery, Iceberg sinks |
| Progress tracked as | source offsets defined by the connector (an LSN, a file position, a timestamp) | Kafka consumer offsets of the group connect-<connector name> |
| Parallelism | tasks split the source (tables, files, shards) | tasks share the topic’s partitions, like consumers in a group |
| Delivery | at-least-once by default; exactly-once possible for connectors that support it (worker setting exactly.once.source.support=enabled, distributed mode only) |
at-least-once by default; exactly-once depends on the connector writing idempotently or transactionally |
Worked example: a sink connector is a consumer group
Each sink connector got its own consumer group named connect-<name>:
$ bin/kafka-consumer-groups.sh --bootstrap-server localhost:19092 --list | grep connect
connect-customers-file-sink-v2
connect-customers-file-sink
connect-raw-events-sink
So everything from the consumers lesson applies: lag, kafka-consumer-groups.sh --describe, and the fact that sink tasks beyond the partition count sit idle.
Pitfalls
- Assuming sinks are exactly-once. Many sinks are at-least-once; check whether yours upserts by key or deduplicates.
- Too few partitions on topics read by a sink:
tasks.max=10on a three-partition topic still runs three busy tasks.
In interviews
Define both, say where each stores its progress, and give a typical pipeline: Debezium source from Postgres, then an S3 or Iceberg sink to the lake.
Standalone vs distributed
What it is. Connect workers run in one of two modes.
| Standalone | Distributed | |
|---|---|---|
| Processes | one | one or more, forming a cluster by sharing group.id |
| Connector configs | properties files on the command line | submitted through the REST API, stored in config.storage.topic |
| Source offsets | a local file (offset.storage.file.filename) |
the compacted Kafka topic offset.storage.topic |
| Fault tolerance | none: the process is a single point of failure | tasks rebalance to surviving workers |
| Exactly-once source support | not available | available |
| Use | local testing, simple edge agents | production |
Worked example: standalone mode
A file source in standalone mode copies a two-line file into a topic:
# file-source.properties
name=lines-source
connector.class=FileStreamSource
tasks.max=1
file=/data/input.txt
topic=standalone-lines
bin/connect-standalone.sh standalone.properties file-source.properties
bin/kafka-console-consumer.sh --bootstrap-server localhost:19092 --topic standalone-lines --from-beginning
line one
line two
The source offset went into the local offsets file: the key ["lines-source",{"filename":".../input.txt"}] and the value {"position":18}, the byte position after the two lines. Lose that file and the connector reads the file again from the start.
Pitfalls
- Running production in standalone mode because it is “simpler”. When that one process dies, nothing takes over.
- Distributed workers with different
group.ids form separate clusters that do not share work. - Replication factor 1 on the internal topics (
connect-configs,connect-offsets,connect-status): losing one broker loses connector state. Use 3.
In interviews
Know where configs and offsets live in each mode, and that distributed workers rebalance connectors and tasks among themselves through the same group protocol idea as consumers.
Converters and SMTs
What it is. A converter controls the bytes in Kafka: JsonConverter, StringConverter, ByteArrayConverter, and the Avro, Protobuf or JSON Schema converters that use a Schema Registry. A single message transform (SMT) is a lightweight per-record transformation: rename or drop a field, add a static field, route to another topic, mask a value, filter records.
How they fit together
Source: external system -> connector -> SMTs -> converter -> Kafka
Sink: Kafka -> converter -> SMTs -> connector -> external system
Converters are set on the worker and can be overridden per connector (key.converter, value.converter). Built-in SMTs include InsertField, ReplaceField, MaskField, Cast, Flatten, RegexRouter, ExtractField, HoistField, TimestampRouter, Filter (with predicates) and header transforms. Debezium adds its own, notably ExtractNewRecordState, which unwraps a change event into the plain row.
Worked example: unwrapping change events into rows
A file sink reads the Debezium topic, unwraps each change event and adds a static source_system field. The Debezium source and this sink both use JsonConverter with schemas enabled:
{
"name": "customers-file-sink-v2",
"config": {
"connector.class": "org.apache.kafka.connect.file.FileStreamSinkConnector",
"topics": "shopv2.public.customers",
"file": "/data/customers-v2.jsonl",
"key.converter": "org.apache.kafka.connect.json.JsonConverter",
"key.converter.schemas.enable": "true",
"value.converter": "org.apache.kafka.connect.json.JsonConverter",
"value.converter.schemas.enable": "true",
"transforms": "unwrap,addSource",
"transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState",
"transforms.unwrap.delete.tombstone.handling.mode": "rewrite",
"transforms.addSource.type": "org.apache.kafka.connect.transforms.InsertField$Value",
"transforms.addSource.static.field": "source_system",
"transforms.addSource.static.value": "shop-db"
}
}
After a snapshot of rows 1 and 3, an update of row 3 and a delete of row 1, the sink file contained flat rows, with deletes rewritten as __deleted=true:
Struct{id=1,[email protected],tier=gold,__deleted=false,source_system=shop-db}
Struct{id=3,[email protected],tier=bronze,__deleted=false,source_system=shop-db}
Struct{id=3,[email protected],tier=silver,__deleted=false,source_system=shop-db}
Struct{id=1,[email protected],tier=gold,__deleted=true,source_system=shop-db}
The schema trap
The first attempt used the worker’s default schemas.enable=false. InsertField still worked, but ExtractNewRecordState silently passed the records through unchanged, because without a schema the value arrives as a plain map rather than a structured record:
{op=u, before=null, source_system=shop-db, ..., after={tier=gold, id=1, [email protected]}, source={...}, ...}
Many SMTs need schema’d data. In production, the usual answer is a schema-aware format (Avro, Protobuf or JSON Schema with a Schema Registry) rather than plain JSON with an embedded schema, which repeats the schema in every message.
Pitfalls
- Converter mismatch. A sink must use the converter that matches how the topic was written; JSON with
schemas.enable=trueexpects a{"schema":...,"payload":...}envelope and fails on plain JSON. - Doing real logic in SMTs. Joins, aggregations and lookups do not belong there; use Kafka Streams, Flink or Spark.
- Ordering of the
transformslist matters: transforms apply left to right.
In interviews
Explain converter versus SMT (format versus content), list a few built-in SMTs, and mention that SMTs are for simple stateless changes. The schemas.enable story is a good example of a real debugging experience.
Offset storage
What it is. Connect must remember how far each connector got. Sink connectors use ordinary consumer group offsets. Source connectors store connector-defined offsets in the offset.storage.topic (distributed) or a local file (standalone).
Worked example: reading offsets
Debezium’s source offset is a position in the PostgreSQL write-ahead log (WAL), the log sequence number (LSN):
$ curl -s localhost:8083/connectors/shop-cdc/offsets
{"offsets":[{"partition":{"server":"shop"},"offset":{"lsn_proc":26625008,"lsn_events_processed":1,"messageType":"DELETE","lsn_commit":26343224,"lsn":26625008,"txId":750,"ts_usec":1791513187968997}}]}
$ curl -s localhost:8083/connectors/raw-events-sink/offsets
{"offsets":[{"partition":{"kafka_partition":0,"kafka_topic":"raw-events"},"offset":{"kafka_offset":3}}]}
The same source offsets are visible as records in the compacted connect-offsets topic, keyed by connector name and source partition:
["shop-cdc",{"server":"shop"}] {"lsn_proc":26625008,"lsn_events_processed":1,"messageType":"DELETE","lsn_commit":26343224,"lsn":26625008,"txId":750,"ts_usec":1791513187968997}
The sink’s offset query returned an empty list for about a minute after it started: workers commit offsets every offset.flush.interval.ms (60 seconds by default), not after every record.
Resetting or editing offsets
Since Kafka 3.6 the REST API can change offsets directly, with the connector stopped:
curl -s -X PUT localhost:8083/connectors/raw-events-sink/stop
curl -s -X PATCH -H 'Content-Type: application/json' localhost:8083/connectors/raw-events-sink/offsets \
-d '{"offsets":[{"partition":{"kafka_topic":"raw-events","kafka_partition":0},"offset":{"kafka_offset":0}}]}'
curl -s -X PUT localhost:8083/connectors/raw-events-sink/resume
DELETE /connectors/{name}/offsets resets them entirely, so a sink restarts from auto.offset.reset and a source starts as if new (for Debezium, that can mean a new snapshot).
Pitfalls
- Deleting and recreating a connector with the same name reuses its old offsets. Use a new name or reset offsets if you want a fresh start.
- Up to a minute of reprocessing after a crash, because offsets are flushed periodically. Sinks must tolerate duplicates.
- Source offsets pointing at data the source no longer has (for example WAL removed): see Debezium’s
snapshot.mode=when_needed.
In interviews
Say where each kind of offset lives, how to read and reset them (REST offsets API, consumer group tools for sinks), and why duplicates after a restart are normal.
Dead letter queue
What it is. By default a connector fails fast: the first record it cannot convert, transform or write stops the task (errors.tolerance=none). With errors.tolerance=all, bad records are skipped, and with errors.deadletterqueue.topic.name they are copied to a dead letter queue topic for later inspection. The DLQ is a sink connector feature.
Worked example
Three records were written to raw-events; the middle one is not JSON:
printf '{"event":"click","page":"/home"}\nnot json at all\n{"event":"click","page":"/cart"}\n' | \
bin/kafka-console-producer.sh --bootstrap-server localhost:19092 --topic raw-events
{
"name": "raw-events-sink",
"config": {
"connector.class": "org.apache.kafka.connect.file.FileStreamSinkConnector",
"topics": "raw-events",
"file": "/data/raw-events.jsonl",
"value.converter": "org.apache.kafka.connect.json.JsonConverter",
"value.converter.schemas.enable": "false",
"errors.tolerance": "all",
"errors.deadletterqueue.topic.name": "raw-events.dlq",
"errors.deadletterqueue.topic.replication.factor": "3",
"errors.deadletterqueue.context.headers.enable": "true",
"errors.log.enable": "true"
}
}
The two good records reached the file, and the bad one landed in raw-events.dlq with headers explaining exactly where and why it failed:
{page=/home, event=click}
{page=/cart, event=click}
__connect.errors.topic:raw-events
__connect.errors.partition:0
__connect.errors.offset:1
__connect.errors.stage:VALUE_CONVERTER
__connect.errors.class.name:org.apache.kafka.connect.json.JsonConverter
__connect.errors.exception.class.name:org.apache.kafka.connect.errors.DataException
__connect.errors.exception.message:Converting byte[] to Kafka Connect data failed due to serialization error:
The DLQ record keeps the original bytes (not json at all), so it can be fixed and replayed.
Pitfalls
errors.tolerance=allwithout a DLQ or logging silently drops data.- Nobody watches the DLQ. Alert on its growth, and have a replay process.
- Retries versus tolerance.
errors.retry.timeoutretries transient failures; tolerance decides what happens to records that still fail. Errors thrown by the external system on write are handled by the connector itself and may not reach the DLQ; check your connector’s behaviour.
In interviews
Describe fail-fast versus tolerate, the DLQ topic with context headers, and the operational loop: alert, inspect, fix, replay. Mention that source connectors do not use the DLQ.
Debezium CDC
What it is. Change data capture turns every insert, update and delete in a database into an event. Debezium does it by reading the database’s own replication log (the PostgreSQL WAL, the MySQL binlog, and so on), so it sees every committed change, including deletes, with low overhead and without polling queries. Debezium runs as Kafka Connect source connectors (it can also run as Debezium Server or an embedded library).
How the PostgreSQL connector works
- PostgreSQL must run with
wal_level = logical. - The connector creates a replication slot (
slot.name) and a publication for the captured tables, and decodes changes with the built-inpgoutputplug-in (plugin.name=pgoutput; managed services such as RDS and Cloud SQL require it). - On first start it takes a snapshot of existing rows (
snapshot.mode=initial, the default), emittingop: "r"events, then streams changes from the slot. - Each table goes to a topic named
<topic.prefix>.<schema>.<table>, keyed by the table’s primary key. - Its offset is the last processed LSN; the slot lets PostgreSQL know which WAL it may discard.
Worked example: capturing changes
{
"name": "shop-cdc",
"config": {
"connector.class": "io.debezium.connector.postgresql.PostgresConnector",
"database.hostname": "localhost",
"database.port": "5499",
"database.user": "postgres",
"database.password": "postgres",
"database.dbname": "shop",
"topic.prefix": "shop",
"plugin.name": "pgoutput",
"slot.name": "shop_cdc",
"publication.autocreate.mode": "filtered",
"table.include.list": "public.customers",
"snapshot.mode": "initial",
"tasks.max": "1"
}
}
curl -s -X POST -H 'Content-Type: application/json' --data @shop-cdc.json localhost:8083/connectors
curl -s localhost:8083/connectors/shop-cdc/status
{"name":"shop-cdc","connector":{"state":"RUNNING","worker_id":"localhost:8083","version":"3.7.0.Final"},"tasks":[{"id":0,"state":"RUNNING","worker_id":"localhost:8083","version":"3.7.0.Final"}],"type":"source"}
Then three changes in one transaction:
UPDATE customers SET tier = 'gold' WHERE id = 1;
INSERT INTO customers VALUES (3, '[email protected]', 'bronze');
DELETE FROM customers WHERE id = 2;
The topic shop.public.customers received six records (values trimmed to before, after and op; each also carries a source block with the LSN, transaction ID, table and timestamps):
Offset:0 {"id":1} {"before":null,"after":{"id":1,"email":"[email protected]","tier":"bronze"},"op":"r",...}
Offset:1 {"id":2} {"before":null,"after":{"id":2,"email":"[email protected]","tier":"silver"},"op":"r",...}
Offset:2 {"id":1} {"before":null,"after":{"id":1,"email":"[email protected]","tier":"gold"},"op":"u",...}
Offset:3 {"id":3} {"before":null,"after":{"id":3,"email":"[email protected]","tier":"bronze"},"op":"c",...}
Offset:4 {"id":2} {"before":{"id":2,"email":"","tier":""},"after":null,"op":"d",...}
Offset:5 {"id":2} null
Reading it:
opisr(snapshot read),c(create),u(update) ord(delete);tmarks a truncate when enabled.- The update’s
beforeisnulland the delete’sbeforecontains only the key (with empty strings for the otherNOT NULLcolumns). That is PostgreSQL’s defaultREPLICA IDENTITY, which logs only the primary key of the old row. AfterALTER TABLE customers REPLICA IDENTITY FULL, the second connector run received full before-images, which is why its delete above carried the email and tier. - The delete is followed by a tombstone (key
{"id":2}, valuenull) so that log compaction can remove the key entirely (see tombstones).
Worked example: applying change events to a replica
A consumer can rebuild the table by upserting after for r, c and u, deleting by key for d, and skipping tombstones. Because the operations are idempotent, replaying the topic does not change the result:
import json, sqlite3
# Key and value of each record on shop.public.customers, trimmed to the fields we use
# (the real values also carry "source" metadata and timestamps).
records = [
('{"id":1}', '{"before":null,"after":{"id":1,"email":"[email protected]","tier":"bronze"},"op":"r"}'),
('{"id":2}', '{"before":null,"after":{"id":2,"email":"[email protected]","tier":"silver"},"op":"r"}'),
('{"id":1}', '{"before":null,"after":{"id":1,"email":"[email protected]","tier":"gold"},"op":"u"}'),
('{"id":3}', '{"before":null,"after":{"id":3,"email":"[email protected]","tier":"bronze"},"op":"c"}'),
('{"id":2}', '{"before":{"id":2,"email":"","tier":""},"after":null,"op":"d"}'),
('{"id":2}', None), # tombstone that follows a delete
]
replica = sqlite3.connect(":memory:")
replica.execute("CREATE TABLE customers (id INTEGER PRIMARY KEY, email TEXT, tier TEXT)")
def apply(key_json, value_json):
key = json.loads(key_json)
if value_json is None:
return "tombstone: skip (it only exists for log compaction)"
event = json.loads(value_json)
if event["op"] in ("r", "c", "u"): # snapshot read, create, update
row = event["after"]
replica.execute("INSERT INTO customers VALUES (:id, :email, :tier) "
"ON CONFLICT (id) DO UPDATE SET email = excluded.email, tier = excluded.tier", row)
return f"upsert id={row['id']}"
if event["op"] == "d":
replica.execute("DELETE FROM customers WHERE id = ?", (key["id"],))
return f"delete id={key['id']}"
for k, v in records:
print(apply(k, v))
for _ in range(2): # replaying the whole topic changes nothing
for k, v in records:
apply(k, v)
print(replica.execute("SELECT * FROM customers ORDER BY id").fetchall())
upsert id=1
upsert id=2
upsert id=1
upsert id=3
delete id=2
tombstone: skip (it only exists for log compaction)
[(1, '[email protected]', 'gold'), (3, '[email protected]', 'bronze')]
The replica matches the source table. Applying events in order per key matters, which is why Debezium keys by primary key: all changes to one row land in one partition.
Snapshot modes
snapshot.mode |
Behaviour |
|---|---|
initial (default) |
snapshot when no offsets exist, then stream |
always |
snapshot on every start, then stream |
initial_only |
snapshot and stop, no streaming |
no_data |
never snapshot; stream from the stored LSN or from the slot’s creation point |
when_needed |
snapshot if there are no offsets, or if the stored offset is no longer available on the server |
For re-snapshotting a table while streaming continues, Debezium supports incremental snapshots triggered through a signalling table or topic.
Pitfalls
- WAL pile-up. A replication slot keeps WAL until the connector confirms it. If the connector is stopped for days, the database disk fills. Monitor slot lag (
pg_replication_slots), and drop slots of retired connectors. - Quiet tables, busy database. If captured tables rarely change while other tables do, the slot’s position may not advance; Debezium’s
heartbeat.interval.ms(optionally with a heartbeat action query) keeps it moving. - Schema changes. Adding a column flows through, but consumers and schemas must handle it; renaming or retyping columns needs a plan (see Schema Registry).
- Missing before-images. Updates and deletes carry only the key unless the table uses
REPLICA IDENTITY FULL. - One slot per connector. Two connectors cannot share
slot.name. - Snapshots are not free. The initial snapshot of a large table reads it all; plan for load and time.
In interviews
“Design CDC from Postgres to the lake.” Cover logical decoding and the replication slot, the initial snapshot, the event envelope (before, after, op, source), keying by primary key for per-row order, tombstones and compaction, idempotent upserts or MERGE in the sink, schema evolution, and the WAL retention risk. The CDC patterns lesson covers the wider design choices.
Practice questions
A sink connector shows RUNNING, but no data arrives in the target. Where do you look?
Check task states, not just the connector state: GET /connectors/{name}/status shows failed tasks with stack traces. Then check the sink’s consumer group (connect-<name>) for lag and whether it is assigned, converter settings that do not match the topic’s format, and whether errors.tolerance=all is silently sending everything to a DLQ.
What is the difference between a converter and an SMT?
A converter controls serialisation, turning Connect’s internal records into bytes in Kafka and back (JSON, Avro, Protobuf). An SMT changes the content of each record (rename, drop or add a field, mask, route, filter) between the connector and the converter. Converters are about format; SMTs are about simple per-record content changes.
Your Debezium connector was stopped over a long weekend and the database disk is nearly full. Why, and what do you do?
The replication slot kept every WAL segment since the connector’s last confirmed LSN, so PostgreSQL could not recycle them. Restart the connector so it catches up and confirms progress, or, if it is retired, drop the slot. Prevent it with alerts on slot lag and disk usage, heartbeats for quiet tables, and a runbook for long outages.
Why does Debezium emit a tombstone after a delete event?
Topics fed by Debezium are often compacted to keep the latest state per primary key. The delete event itself has a value, so compaction would keep it; the following tombstone (same key, null value) lets compaction remove the key entirely after delete.retention.ms. Consumers that apply changes should treat the delete event as the delete and skip the tombstone.
How does a Connect cluster tolerate a worker failure in distributed mode?
Workers sharing a group.id coordinate through Kafka. When one leaves or dies, the remaining workers rebalance its connectors and tasks among themselves; configs come from the config topic, source offsets from the offset topic and sink offsets from consumer groups, so tasks resume from their last committed positions (with some reprocessing). Standalone mode has none of this.
How would you replay the last day of data into a sink connector?
Stop the connector (PUT /connectors/{name}/stop), use PATCH /connectors/{name}/offsets to set each partition’s kafka_offset (or reset the connect-<name> group with kafka-consumer-groups.sh --reset-offsets --by-duration P1D), then resume. Make sure the sink writes idempotently, otherwise the replay creates duplicates.
Key takeaways
- Kafka Connect runs reusable connectors as workers, connectors and tasks, configured through a REST API instead of code.
- Source connectors bring data into Kafka and track their own offsets; sink connectors are consumer groups named
connect-<name>. - Use distributed mode in production: configs, offsets and status live in replicated Kafka topics, and tasks fail over between workers.
- Converters set the format, SMTs make small per-record changes, and many SMTs need schema’d data.
- Default error handling fails fast;
errors.tolerance=allwith a DLQ and context headers keeps a sink running while keeping bad records. - Debezium reads the database log through a replication slot, snapshots first, keys events by primary key, and follows deletes with tombstones; watch WAL retention and replica identity.
Progress is saved in this browser only. No account needed.

