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 pageshowhide
explore
- Connect Framework and Runtime6 questions
- Connect REST API and Lifecycle Management6 questions
- Source vs Sink Connectors and Tasks5 questions
- Converters and Schema Handling6 questions
- Single Message Transforms and Predicates5 questions
- Offset Management in Connect5 questions
- Error Handling and Dead-Letter Queue5 questions
- Change Data Capture with Debezium6 questions
- Scaling and Rebalancing Across Workers6 questions
- Connect Security and Config Externalization5 questions
questions
page 2 of 2What do errors.log.enable and errors.log.include.messages do, and what is the security consideration with the latter?
basics
~10 serrors.log.enable=true logs each failed record's error context to the Connect worker log. errors.log.include.messages=true additionally logs the record's key/value contents. Including messages can leak sensitive payload data into logs.
What are the essential worker-level configuration properties needed to bootstrap a distributed Connect worker, and what does each control?
basics
~20 sYou need bootstrap.servers (Kafka brokers), group.id (cluster identity), the three internal topic names with replication factors, key/value converters, plugin.path, and the REST listener. These let the worker connect to Kafka, join its cluster, and load plugins.
What is the OffsetStorageReader, and how should a source task use it on startup?
basics
~20 sOffsetStorageReader lets a source task read back the last committed source offsets when it starts, so it can resume from where it left off. The task gets it from its SourceTaskContext and queries by source partition.
Walk through pausing, stopping, and resuming a connector via the REST API. How do PUT /pause, PUT /stop, and PUT /resume differ?
basics
~10 sPUT /connectors/{name}/pause halts processing but keeps tasks assigned (state PAUSED). PUT /connectors/{name}/stop also halts but de-allocates the tasks (state STOPPED), freeing cluster resources. PUT /connectors/{name}/resume restarts a paused or stopped connector back to RUNNING.
How does tasks.max interact with worker count when balancing load, and what are the failure modes of misconfiguring it?
basics
~20 stasks.max caps how many tasks a connector creates, which caps how many workers can share its load. Too low wastes workers (idle) and limits throughput; too high creates more tasks than the source can usefully split, adding overhead with no gain.
How does a sink connector decide which Kafka topics to consume, and how are partitions assigned to its tasks?
basics
~20 sA sink connector subscribes to topics via the topics config (explicit list) or topics.regex (pattern). The tasks form a consumer group, so Kafka's group rebalance assigns each partition of those topics to exactly one task.
Explain how the Postgres Debezium connector consumes the WAL: logical replication, output plugins (pgoutput/wal2json), replication slots, and the operational risks of slots.
basics
~20 sThe Postgres connector uses logical replication: Postgres must run with wal_level=logical. A logical decoding output plugin (pgoutput, built in, or wal2json) turns WAL records into change events, delivered through a replication slot that tracks the last LSN the connector confirmed. The big risk: if the connector is down, the slot holds WAL on disk and can fill the disk.
A sink connector throws DataException / SerializationException on deserialize. How do you diagnose and resolve converter and schema mismatches?
basics
~20 sThe sink's converter usually doesn't match how the data was actually written. Identify the real on-topic format, align the sink's key/value converter and settings (e.g. schemas.enable, schema.registry.url) to it, and use error tolerance / a DLQ for poison records.
Which categories of failures does Kafka Connect's error-handling framework tolerate, and which does it NOT? Distinguish conversion/transform/put failures from other error types.
basics
~20 sThe framework tolerates errors in the parts it controls: converters (deserialization/serialization), SMTs, and the sink connector's put() handoff. It does NOT tolerate failures outside that boundary — like consumer/producer/offset-commit errors or framework/configuration errors, which still fail the task.
Explain errors.retry.timeout and errors.retry.delay.max.ms. How do retries interact with errors.tolerance?
basics
~20 serrors.retry.timeout is how long Connect keeps retrying a failed operation (default 0 = no retries; -1 = forever). errors.retry.delay.max.ms caps the backoff between attempts (default 60000). Retries run first; only after they're exhausted does errors.tolerance decide skip vs fail.
What is plugin.path in Kafka Connect, and how does the runtime isolate connector plugins? What problems does misconfiguring it cause?
basics
~10 splugin.path is a list of directories where the worker finds connector/converter/transform JARs. Connect loads each plugin in an isolated classloader to avoid dependency conflicts. Misconfiguring it causes 'class not found' errors or dependency clashes.
How does exactly-once delivery for SOURCE connectors work in Kafka Connect, and what is required to enable it?
basics
~20 sConnect (KIP-618, Kafka 3.3+) can run source connectors exactly-once by writing the produced records and their source offsets together in a single Kafka transaction. You enable it cluster-wide with exactly.once.source.support=enabled, and the connector must declare exactly.once.support and define transaction boundaries.
A team needs to make a source connector reprocess data from the beginning, and separately rewind a sink connector. How do you reset offsets for each, safely?
basics
~20 sStop the connector first. For modern Connect (3.6+), use the REST offsets API: GET/DELETE/PATCH /connectors/{name}/offsets. DELETE wipes offsets so a source restarts from scratch; PATCH sets specific source offsets or sink consumer-group offsets. Older clusters edit offset.storage.topic (source) or use kafka-consumer-groups --reset-offsets (sink).
How do you keep secrets like passwords out of a connector config submitted to the REST API? Explain ConfigProvider externalization.
basics
~20 sUse a ConfigProvider: put a placeholder like ${file:/path:key} or ${vault:...} in the connector config instead of the literal secret. Connect resolves it at runtime from the provider, so the secret never lives in the REST payload or the config topic, and GET /config shows the placeholder.
How does Kafka Connect validate a connector configuration before accepting it, and how does PUT /connector-plugins/{class}/config/validate help you build safe deployment tooling?
basics
~20 sConnect runs each connector's Validator over the proposed config. PUT /connector-plugins/{class}/config/validate returns a structured per-field result (definitions, current values, recommended values, and per-field errors) without creating anything, so tooling can pre-flight a config and surface errors before deploying.
Describe the connector restart endpoint and the includeTasks and onlyFailed query parameters. How would you restart just the failed tasks of a connector?
basics
~10 sPOST /connectors/{name}/restart restarts only the Connector instance by default. Add ?includeTasks=true to also restart tasks, and ?onlyFailed=true to limit the restart to FAILED instances. To restart just failed tasks: POST /connectors/{name}/restart?includeTasks=true&onlyFailed=true.
What does scheduled.rebalance.max.delay.ms control, and how do you tune it for worker restarts vs. genuine failures?
basics
~20 sIt's how long (default 5 minutes) the leader waits before reassigning a departed worker's tasks, hoping the worker returns and reclaims them. It avoids churn from quick restarts but keeps those tasks down during the wait if the worker is really dead.
How do you secure the Kafka Connect REST API with TLS and authentication, and what's the listeners.https configuration?
basics
~10 sSet listeners=https://host:port plus listeners.https.ssl.* keystore/truststore configs for TLS. For auth, plug in a REST extension via rest.extension.classes (e.g. BasicAuthSecurityRestExtension for HTTP Basic) since Connect has no built-in user store.
How do you implement a custom SMT by implementing the Transformation interface, and what are the key methods and gotchas (schema handling, key/value variants, deployment)?
basics
~20 sImplement org.apache.kafka.connect.transforms.Transformation<R extends ConnectRecord<R>>. Override apply(R) to return a transformed record, configure(Map) to read config, config() to declare a ConfigDef, and close() to release resources. Handle both schema and schemaless records, build a new record with record.newRecord(...), package it as a plugin, and put the Jar on the Connect plugin path.
Compare SourceRecord and SinkRecord. What fields does each carry, and why does SourceRecord have sourcePartition/sourceOffset while SinkRecord has kafkaPartition/kafkaOffset?
basics
~20 sSourceRecord describes data going INTO Kafka: target topic/partition, key, value, schemas, plus sourcePartition/sourceOffset that point back into the external system for resume. SinkRecord describes data coming OUT of Kafka: it carries the actual Kafka topic, kafkaPartition, and kafkaOffset of where it was consumed.
What delivery guarantee does Debezium provide, how do you achieve effectively-exactly-once end to end, and what is the role of exactly-once snapshot/streaming and Kafka Connect EOS support?
basics
~20 sDebezium delivers at-least-once: after a crash it may re-emit some events because offsets are committed periodically, not per-event. Effective exactly-once is achieved downstream by idempotency — each event carries the primary key, so consumers/sinks upsert by key and re-applied duplicates are harmless. Kafka Connect 3.3+ added exactly-once source support that can make source-side delivery exactly-once.
How do workers in a distributed Connect cluster coordinate work assignment, and how did incremental cooperative rebalancing (KIP-415) improve the original protocol?
basics
~20 sWorkers sharing a group.id use Kafka's group-membership protocol to elect a leader that assigns connectors and tasks. The original protocol stopped all work on every change (stop-the-world); incremental cooperative rebalancing (KIP-415) only reassigns the affected tasks, avoiding global pauses.
How would you design a rolling restart / upgrade of a large Connect cluster to minimize rebalancing disruption?
basics
~20 sUse cooperative rebalancing (connect.protocol=compatible/sessioned) so only moving tasks pause. Set scheduled.rebalance.max.delay.ms above your per-worker restart time so restarting workers reclaim their own tasks. Restart one worker at a time and wait for the cluster to settle between each.
Design how a multi-tenant Connect cluster gives each connector its own broker identity for per-tenant ACLs. What pieces fit together?
basics
~10 sSet connector.client.config.override.policy=Principal on the worker, then give each connector its own credentials via producer.override./consumer.override.sasl.jaas.config (resolved from a ConfigProvider). Define per-principal broker ACLs so each tenant can only touch its own topics.
When designing an SMT chain (e.g. ExtractField + Cast + Flatten + RegexRouter with predicates), why does ordering matter and what are common pitfalls?
basics
~20 sSMTs run in listed order and each consumes the previous one's output, so an SMT that renames or removes a field changes what later SMTs can see. You must order field-shaping before routing or extraction that depends on those fields, mind schema changes, predicate scoping, and tombstone handling.
showing 31–55 of 55