Menu

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.

  • Intermediate
  • 20 min read
  • Updated Oct 2026
On this page
  1. Sample setup
  2. Connect framework
  3. The pieces
  4. Worked example: the REST API
  5. Pitfalls
  6. In interviews
  7. Source vs sink connectors
  8. Worked example: a sink connector is a consumer group
  9. Pitfalls
  10. In interviews
  11. Standalone vs distributed
  12. Worked example: standalone mode
  13. Pitfalls
  14. In interviews
  15. Converters and SMTs
  16. How they fit together
  17. Worked example: unwrapping change events into rows
  18. The schema trap
  19. Pitfalls
  20. In interviews
  21. Offset storage
  22. Worked example: reading offsets
  23. Resetting or editing offsets
  24. Pitfalls
  25. In interviews
  26. Dead letter queue
  27. Worked example
  28. Pitfalls
  29. In interviews
  30. Debezium CDC
  31. How the PostgreSQL connector works
  32. Worked example: capturing changes
  33. Worked example: applying change events to a replica
  34. Snapshot modes
  35. Pitfalls
  36. In interviews
  37. Practice questions
  38. 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 with wal_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 RUNNING connector with a FAILED task. 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=10 on 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=true expects 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 transforms list 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=all without 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.timeout retries 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

  1. PostgreSQL must run with wal_level = logical.
  2. The connector creates a replication slot (slot.name) and a publication for the captured tables, and decodes changes with the built-in pgoutput plug-in (plugin.name=pgoutput; managed services such as RDS and Cloud SQL require it).
  3. On first start it takes a snapshot of existing rows (snapshot.mode=initial, the default), emitting op: "r" events, then streams changes from the slot.
  4. Each table goes to a topic named <topic.prefix>.<schema>.<table>, keyed by the table’s primary key.
  5. 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:

  • op is r (snapshot read), c (create), u (update) or d (delete); t marks a truncate when enabled.
  • The update’s before is null and the delete’s before contains only the key (with empty strings for the other NOT NULL columns). That is PostgreSQL’s default REPLICA IDENTITY, which logs only the primary key of the old row. After ALTER 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}, value null) 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=all with 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.

By DataDank Editorial · Last reviewed Oct 2026 · Kafka Connect 4.3.1 (distributed and standalone workers) ran against a three-broker Apache Kafka 4.3.1 KRaft cluster with the bundled FileStream connectors, and Debezium 3.7.0.Final (PostgreSQL connector, pgoutput) captured changes from a separate PostgreSQL 16 instance with wal_level=logical. REST calls, connector output, DLQ headers and offsets on this page are from those runs. The replica-table example runs on plain Python 3 with SQLite.

Progress is saved in this browser only. No account needed.

Search
Filter by type