Kafka Connect
Kafka Connect is Kafka's framework for streaming data between Kafka and other systems (databases, object storage, search engines, SaaS APIs) using configurable connectors instead of hand-written producer and consumer code. A source connector pulls data from an external system into Kafka topics; a sink connector pushes data from topics into an external system. Connect handles the hard parts: parallelism, offset tracking, fault tolerance, retries, and schema conversion.
It's the backbone of many data platforms: Debezium streams database changes into Kafka, and sink connectors load them into Elasticsearch, S3, Snowflake, or a data lakehouse. Hundreds of connectors exist, so the usual job is choosing, configuring, and operating them well.
TL;DR
- Source connectors bring data into Kafka; sink connectors send data out of Kafka.
- Connectors run in workers; each connector is split into tasks for parallelism.
- Run in distributed mode in production: workers form a cluster, share load, and store config, offsets, and status in Kafka topics.
- Converters (Avro, Protobuf, JSON Schema with Schema Registry) control the data format on the topic.
- Single Message Transforms (SMTs) make lightweight per-record changes (rename fields, route topics, mask data).
- Configure error tolerance and dead-letter queues so bad records don't stop a pipeline.
Quick Example
A Debezium PostgreSQL source and an S3 sink, created via the Connect REST API:
No application code: two JSON configs create a CDC pipeline from Postgres to a data lake.
Core Concepts
Connectors, Tasks, and Workers
- A connector is a configured instance of a connector plugin (for example "the S3 sink for orders").
- The connector splits its work into tasks, up to
tasks.max. For sinks, tasks divide topic partitions; for sources, they divide tables, files, or other units. - Workers are the JVM processes that run tasks. In distributed mode, workers with the same
group.idform a cluster: tasks are balanced across workers and rebalanced automatically when a worker fails. - Distributed mode stores configurations, source offsets, and status in internal Kafka topics, so workers are stateless and easy to run on Kubernetes (for example with the Strimzi operator).
Standalone mode runs one worker with local config files, which is fine for development and edge cases, but it has no fault tolerance.
Converters and Schemas
Connectors work with Connect's internal data model; converters serialize it to bytes on the topic:
Sink connectors that write to typed systems (databases, Parquet, warehouses) need schemas to create tables and map types correctly.
Single Message Transforms
SMTs apply simple, stateless changes to each record as it flows through a connector:
ReplaceField,MaskField,Cast,InsertField: shape and sanitize data.RegexRouter,TimestampRouter: change destination topic names.ExtractField,ValueToKey,HoistField: restructure keys and values.- Debezium's
ExtractNewRecordState: flatten change events to the "after" state. - Predicates apply transforms conditionally.
Keep SMTs lightweight. Joins, aggregations, and enrichment belong in a stream processing layer (Kafka Streams, Flink).
Change Data Capture With Debezium
Debezium source connectors read a database's replication log (Postgres logical decoding, MySQL binlog, SQL Server CDC, MongoDB change streams) and emit an event for every insert, update, and delete, with before and after states, plus an initial snapshot. It's the standard way to feed databases into Kafka without dual writes, and the relay for the transactional outbox pattern. See change data capture.
Error Handling and Delivery
errors.tolerance=none(default): any conversion or transform error fails the task, which is safe but stops the pipeline.errors.tolerance=allwitherrors.deadletterqueue.topic.name(sinks): bad records go to a DLQ topic with error context in headers, and the pipeline continues.errors.retry.timeoutanderrors.retry.delay.max.msretry transient failures.errors.log.enable=truelogs failing records for debugging.
Sink delivery is typically at-least-once. Many sinks achieve effectively-once results with idempotent writes (upserts keyed by primary key, deterministic object names in S3). Source connectors can support exactly-once in distributed mode on Kafka 3.3+ (KIP-618) when the connector implements it.
Best Practices
Run Distributed Connect as a Platform
Treat Connect clusters as shared infrastructure, or run one per domain or team for isolation. Manage connector configs as code (GitOps, Strimzi KafkaConnector resources, or Terraform providers) rather than ad hoc REST calls.
Keep Secrets Out of Connector Configs
Connector configs are readable through the REST API. Use config providers (FileConfigProvider, Vault, or cloud secret managers) so configs contain references, not credentials. See secrets management.
Monitor Connector and Task Status
Alert on tasks in FAILED state, sink consumer lag, source lag (such as Debezium's replication slot lag), DLQ volume, and throughput. A failed task doesn't restart itself unless you configure automation.
Watch Database Side Effects of CDC
Postgres replication slots retain WAL until Debezium reads it. A stopped connector can fill the database disk. Monitor slot lag and set max_slot_wal_keep_size. See PostgreSQL logical replication.
Common Mistakes
Using JSON Without Schemas for Typed Sinks
A JDBC or warehouse sink reading schemaless JSON can't infer column types reliably and fails, or creates everything as strings. Use a schema-aware converter.
Doing Heavy Processing in SMTs
Chaining many transforms, or writing custom SMTs that call external services, makes connectors slow and fragile. Move enrichment and joins into a stream processor.
Leaving errors.tolerance=all Without a DLQ
Tolerating errors without a dead-letter topic means bad records are silently skipped: data loss that nobody notices until someone asks why rows are missing.
FAQ
When should I use Kafka Connect instead of writing a producer or consumer?
When a mature connector exists for the external system, which covers most databases, cloud storage, search engines, and warehouses. Connect gives you scaling, offset management, retries, and monitoring for free. Write custom code when the integration involves complex business logic, or when no connector exists.
What's the difference between Debezium and the JDBC source connector?
The JDBC source connector polls tables with queries (by incrementing ID or timestamp), so it misses deletes and intermediate updates, and it adds query load. Debezium reads the database's transaction log, capturing every insert, update, and delete in order with low overhead. Prefer Debezium for real CDC.
Is Kafka Connect exactly-once?
Source connectors can achieve exactly-once delivery into Kafka in distributed mode (Kafka 3.3+) if the connector supports it. Sink connectors are generally at-least-once. Design sinks to be idempotent (upserts, deterministic file naming) so duplicates don't cause incorrect results.
Do I need Confluent Platform to run Connect?
No. Kafka Connect is part of Apache Kafka. Some connectors are licensed by Confluent or other vendors, while many (Debezium, Apache Camel connectors, and numerous community ones) are open source. Managed Connect services exist on Confluent Cloud, Amazon MSK Connect, Aiven, and others.
Related Topics
- Kafka — The platform overview
- Change Data Capture — Streaming database changes with Debezium
- Kafka Producers — The outbox pattern and reliable publishing
- Stream Processing — Transformations beyond SMTs
- Data Engineering — Pipelines built on Connect
- Enterprise Integration — Integrating systems at scale