skip to content

Kafka Connect

Kafka Connect as the configuration-driven integration layer: workers, source and sink connectors, converters, transforms, dead-letter queues, and CDC. Interviewers ask because Connect is usually the right answer instead of a bespoke producer.

part ofApache Kafkaoverview, primer and where to startread it →
on this pageshow

explore

questions

page 1 of 2

What is Change Data Capture (CDC) with Debezium, and why is log-based CDC preferred over query-based polling?

level: juniorimportance: must knowfreq 78%

answer

  1. reads the transaction log, not the table
  2. MySQL binlog / Postgres WAL
  3. one event per insert/update/delete
  4. polling misses DELETEs
  5. Kafka Connect source connector

basics

~20 s

Debezium is a Kafka Connect source connector that reads a database's transaction log (MySQL binlog, Postgres WAL) and streams every row insert, update, and delete to Kafka as events. Log-based CDC catches all changes, including deletes, with low impact on the database.

solid answer

~40 s

Change Data Capture means capturing row-level changes (inserts/updates/deletes) from a database and emitting them as a stream. Debezium is a set of Kafka Connect source connectors that do this by reading the database's commit log directly: MySQL's binary log (binlog), Postgres's Write-Ahead Log (WAL), etc. It runs inside Connect workers, deserializes log entries, and produces one Kafka message per change. Log-based CDC is preferred over query-based polling (SELECT ... WHERE updated_at > last_run) because the log records every committed change in order — so deletes and intermediate updates are never missed, there is no need for an updated_at column, ordering is preserved, and the database is barely impacted since reading the log is cheap compared to repeated full-table scans.

go deeper

for a junior

Know that Debezium streams row changes from the DB log into Kafka and that it catches deletes which polling misses.

for a middle

Explain the binlog/WAL mechanism and the concrete advantages over polling (deletes, ordering, no updated_at column, low load).

for a senior

Discuss the required DB config (ROW binlog format, wal_level=logical, replication slots) and operational trade-offs of log-based capture.

for a principal

Reason about where CDC fits in an event-driven architecture, downstream contract stability, and when polling or outbox patterns are preferable.

**Change Data Capture (CDC)** is the practice of detecting and capturing changes made to data in a database so other systems can react to them. Instead of periodically asking 'what changed?', CDC delivers a continuous stream of change events. **Debezium** is an open-source CDC platform built on top of **Kafka Connect** (the pluggable integration framework that ships with Apache Kafka). Debezium provides *source connectors* — one per database engine (MySQL, PostgreSQL, MongoDB, SQL Server, Oracle, Db2, Cassandra). A source connector pulls data *into* Kafka. Each Debezium connector runs as one or more *tasks* inside Kafka Connect worker JVMs. **How log-based CDC works:** Relational databases maintain a durable, ordered transaction log to guarantee crash recovery and replication: - **MySQL** writes a **binary log (binlog)** in `ROW` format — every committed row change is recorded. - **PostgreSQL** writes a **Write-Ahead Log (WAL)**; Debezium consumes it via *logical replication* (a logical decoding plugin like `pgoutput` or `wal2json`). Debezium connects as a replication client, reads these log entries as the database commits them, decodes each one into a structured change event, and publishes it to a Kafka topic (one topic per table by default). Because the log is the same mechanism the database uses for its own replicas, Debezium sees an exact, ordered record of every committed change. **Why log-based beats query-based polling:** - **Query-based polling** runs something like `SELECT * FROM t WHERE updated_at > :last`. It requires an `updated_at` column, **misses DELETEs entirely** (a deleted row no longer matches any query), can miss intermediate states (two updates between polls collapse to one), adds query load to the primary, and can lag. - **Log-based CDC** captures *every* committed change including deletes, preserves commit order, needs no extra columns, and imposes minimal load because reading the log is what replicas already do. **Trade-off:** log-based CDC requires database privileges and config (binlog enabled in ROW format with `binlog_row_image=FULL`; in Postgres `wal_level=logical` plus a replication slot and publication). It is operationally heavier to set up but far more correct and efficient at runtime.

  • Why can't query-based polling capture deletes?
    A polling query (WHERE updated_at > last) can only return rows that still exist. A deleted row no longer matches any SELECT, so its disappearance is invisible. The transaction log, by contrast, contains an explicit delete record.
  • Does Debezium run inside the database or separately?
    Separately — it runs as a connector/task inside a Kafka Connect worker process, connecting to the database as a replication client over the network. It is not a database plugin (though Postgres needs a logical-decoding output plugin server-side).

saying these in an interview costs you the question

  • Saying Debezium polls the table with SELECT statements (it reads the commit log).
  • Claiming Debezium runs inside the database server itself.
  • Saying query-based CDC can capture deletes just as well.
  • Confusing Debezium (a source connector) with a sink that writes to a database.

context

open as a page

What is a converter in Kafka Connect, and what is the difference between key.converter and value.converter?

level: juniorimportance: must knowfreq 80%

basics

~10 s

A converter serializes/deserializes Connect data to/from bytes on the Kafka topic. key.converter handles the record key; value.converter handles the record value. They are configured independently and can differ.

open as a page

What does errors.tolerance control in Kafka Connect, and what is the difference between the values none and all?

level: juniorimportance: must knowfreq 78%

basics

~10 s

errors.tolerance decides what Connect does when a record fails processing. none (the default) stops the connector task on the first error; all skips the bad record and keeps going.

open as a page

What is Kafka Connect, and what problem does it solve compared to writing your own producer/consumer applications?

level: juniorimportance: must knowfreq 70%

basics

~10 s

Kafka Connect is a framework for streaming data between Kafka and external systems (databases, files, S3) using ready-made connectors, so you configure plugins instead of writing custom producer/consumer code.

open as a page

How does Kafka Connect track offsets for source connectors versus sink connectors? Where is each kind stored?

level: juniorimportance: must knowfreq 75%

basics

~20 s

Source connectors store their own progress (where they read from the external system) in a special Kafka topic called the offset storage topic. Sink connectors just use a normal Kafka consumer group, so Kafka tracks their offsets like any consumer.

open as a page

How do you create a new connector via the Kafka Connect REST API, and what is the minimal request body?

level: juniorimportance: must knowfreq 75%

basics

~10 s

Send a POST to /connectors with a JSON body containing a top-level "name" and a "config" object (which must repeat connector.class and the connector's settings). Connect returns 201 Created with the connector info.

open as a page

How do you scale a Kafka Connect distributed cluster, and how does work get distributed across workers?

level: juniorimportance: must knowfreq 70%

basics

~10 s

Run multiple Connect worker processes with the same group.id. You scale out by starting more workers. The cluster automatically spreads connectors and their tasks evenly across all available workers via a rebalance.

open as a page

How do you configure a Kafka Connect worker to connect to brokers secured with TLS and SASL, and which prefixes apply to which client?

level: juniorimportance: must knowfreq 70%

basics

~20 s

Put the broker security settings (security.protocol, SASL/SSL configs) in the worker properties file. By default they apply to all of Connect's internal clients. Prefix with producer., consumer., or admin. to target a specific client type.

open as a page

What is a Single Message Transform (SMT) in Kafka Connect, and how do you configure a chain of them on a connector?

level: juniorimportance: must knowfreq 70%

basics

~20 s

An SMT is a small function that modifies each record as it flows through a Kafka Connect connector. You list transforms by alias in the transforms config, then configure each alias with transforms.<alias>.type and its properties. They run in the order listed.

open as a page

What is the difference between a Source connector and a Sink connector in Kafka Connect, and which direction does data flow in each?

level: juniorimportance: must knowfreq 80%

basics

~20 s

A Source connector pulls data from an external system INTO Kafka topics. A Sink connector reads data FROM Kafka topics and writes it OUT to an external system. Source = into Kafka, Sink = out of Kafka.

open as a page

Describe the structure of a Debezium change event envelope (op/before/after/source) and explain how deletes and tombstones work.

level: middleimportance: must knowfreq 72%

basics

~20 s

Each Debezium event has a value with op (c/u/d/r), before (old row), after (new row), and source metadata. A delete sends an event with op='d', before populated, after null — followed by a separate tombstone message (same key, null value) so log-compacted topics can drop the key.

open as a page

Explain JsonConverter and the schemas.enable setting. What changes in the on-wire payload when it is true vs false?

level: middleimportance: must knowfreq 75%

basics

~10 s

JsonConverter serializes Connect data as JSON. With schemas.enable=true it wraps the payload as {"schema":...,"payload":...} so the schema travels inline. With false it emits just the raw JSON value with no schema.

open as a page

How do you configure a dead-letter queue for a Kafka Connect sink connector, and what does context.headers.enable add?

level: middleimportance: must knowfreq 70%

basics

~10 s

Set errors.tolerance=all plus errors.deadletterqueue.topic.name=<topic>. Failed records are written to that topic. errors.deadletterqueue.context.headers.enable=true adds Kafka headers describing why and where the record failed.

open as a page

Compare standalone mode and distributed mode in Kafka Connect. When would you choose each, and how does state/configuration storage differ?

level: middleimportance: must knowfreq 75%

basics

~20 s

Standalone runs a single worker and stores offsets in a local file — simple, no fault tolerance, good for dev. Distributed runs multiple coordinated workers that store config, offsets, and status in Kafka topics, giving scalability and fault tolerance.

open as a page

Walk through how a source task's offsets get committed. What does offset.flush.interval.ms control, and what is the relationship between producing records and committing offsets?

level: middleimportance: must knowfreq 60%

basics

~20 s

A source task hands records to the framework, which produces them to Kafka. Periodically (every offset.flush.interval.ms, default 60000 ms) the framework flushes the producer, then writes the source offsets of acknowledged records to offset.storage.topic. Offsets are only committed after the data is safely in Kafka.

open as a page

Explain the difference between GET /connectors/{name}/status, /config, and /tasks. What does the status endpoint tell you that /config does not?

level: middleimportance: must knowfreq 70%

basics

~20 s

/config returns the desired configuration; /tasks returns the per-task configs; /status returns runtime health: the connector's RUNNING/FAILED/PAUSED state plus each task's state, assigned worker, and any failure trace. Status is live; config is static intent.

open as a page

What problem did the original eager rebalancing protocol in Kafka Connect cause, and why was it painful at scale?

level: middleimportance: must knowfreq 60%

basics

~20 s

The old 'eager' protocol was stop-the-world: on any membership change every worker revoked ALL its connectors and tasks, then everything was reassigned. So one worker joining or a brief blip paused the entire cluster's data flow.

open as a page

How does ConfigProvider (e.g. FileConfigProvider) let you keep secrets out of connector configs, and how does the ${...} indirection work?

level: middleimportance: must knowfreq 65%

basics

~10 s

A ConfigProvider resolves ${provider:[path:]key} placeholders at runtime so secrets aren't stored literally in connector configs. You register a provider in the worker config (e.g. config.providers=file with FileConfigProvider) and reference values like ${file:/secrets.properties:db.password}.

open as a page

What does connector.client.config.override.policy do, and why is it set to None by default?

level: middleimportance: must knowfreq 60%

basics

~10 s

It's a worker-level setting that controls whether individual connectors can override the worker's Kafka client configs using producer.override./consumer.override./admin.override. prefixes. Defaults to None (no overrides allowed); common values are All or Principal.

open as a page

Walk through the common built-in SMTs (InsertField, ReplaceField, MaskField, RegexRouter, TimestampRouter, ExtractField, Cast, Flatten) and what each does.

level: middleimportance: must knowfreq 65%

basics

~20 s

InsertField adds a field; ReplaceField renames/drops fields; MaskField hides values; RegexRouter and TimestampRouter rewrite the target topic name; ExtractField pulls one field up to be the whole key/value; Cast changes field types; Flatten collapses nested structs into dotted top-level fields.

open as a page

Explain the relationship between a Connector and its Tasks. What is the role of taskConfigs() and tasks.max?

level: middleimportance: must knowfreq 72%

basics

~20 s

A Connector is a single coordinator instance that splits work into Tasks, the things that actually move data. taskConfigs(maxTasks) returns the config for each task; tasks.max caps how many tasks Connect will run for that connector.

open as a page

How does Debezium's initial snapshot work, and how does it hand off to streaming the transaction log without missing or duplicating data?

level: seniorimportance: must knowfreq 60%

basics

~20 s

On first start, Debezium reads the current contents of the captured tables (the snapshot) and emits each row as an op='r' event, recording the log position (binlog offset / LSN) at snapshot start. It then begins streaming the transaction log from that recorded position, so every change after the snapshot is captured exactly once.

open as a page

Compare AvroConverter, ProtobufConverter, and JsonConverter, and explain how Schema Registry integration works for the registry-backed converters.

level: seniorimportance: must knowfreq 70%

basics

~20 s

AvroConverter and ProtobufConverter store schemas in Schema Registry and put only a schema ID + binary payload on the topic; JsonConverter embeds or omits the schema inline with no registry. The registry-backed ones give compact messages and enforced compatibility.

open as a page

Describe the three internal topics a distributed Connect cluster uses (config, offset, status). What does each store, and what configuration do they require?

level: seniorimportance: must knowfreq 60%

basics

~10 s

config.storage.topic holds connector/task configs (single partition, compacted). offset.storage.topic holds source-connector read positions (many partitions, compacted). status.storage.topic holds connector/task/worker status (many partitions, compacted). All need high replication.

open as a page

Explain incremental cooperative rebalancing (KIP-415) in Kafka Connect: how does it differ from eager, and what is connect.protocol?

level: seniorimportance: must knowfreq 65%

basics

~20 s

KIP-415 stops revoking everything. Only tasks that must move are revoked; the rest keep running, so there's no cluster-wide pause. It uses two rebalance rounds (revoke, then assign). The connect.protocol setting selects eager, compatible, or sessioned.

open as a page

How do predicates work with SMTs in Kafka Connect, and what do the built-in TopicNameMatches and HasHeaderKey predicates plus `negate` let you do?

level: seniorimportance: must knowfreq 50%

basics

~20 s

A predicate is a per-record boolean test you attach to an SMT so the SMT only runs when the predicate is true. You define predicates under predicates, reference one via transforms.<alias>.predicate, and flip it with transforms.<alias>.negate=true. Built-ins include TopicNameMatches (regex on topic) and HasHeaderKey.

open as a page

Walk through the SourceTask.poll() loop and the SinkTask.put()/flush() loop. What is each responsible for, and how do offsets get committed?

level: seniorimportance: must knowfreq 60%

basics

~20 s

SourceTask.poll() returns new SourceRecords for Connect to produce into Kafka; Connect tracks the source offsets you embed. SinkTask.put() receives batches of SinkRecords to write downstream; flush()/preCommit() is where the sink ensures records are durable before Kafka offsets are committed.

open as a page

What is a Debezium schema change topic (and schema history), and how does Debezium handle DDL/schema evolution in the source database?

level: middleimportance: should knowfreq 42%

basics

~20 s

Debezium tracks the source table structure so it can correctly decode log entries that only carry column positions/values. MySQL persists captured DDL to an internal schema history topic, and can also publish DDL to a separate schema change topic for consumers. When columns are added/changed, the event schema evolves accordingly.

open as a page

Describe Connect's internal data model — Schema and Struct — and how converters relate to it.

level: middleimportance: should knowfreq 55%

basics

~10 s

Connect represents data internally with a Schema (the type/shape) and a Struct (the values), independent of any wire format. Converters translate between this Schema/Struct model and the raw bytes on the topic.

open as a page

When would you use StringConverter, ByteArrayConverter, and a header.converter? What are their constraints?

level: middleimportance: should knowfreq 45%

basics

~10 s

StringConverter treats data as plain text (UTF-8 by default). ByteArrayConverter passes raw bytes through untouched (no schema). header.converter serializes record headers separately, defaulting to SimpleHeaderConverter.

open as a page

showing 1–30 of 55