Learn Labs
9. Building Data Pipelines (Kafka Connect)

9.4 Kafka Connect architecture

Connect parallelism

connector
transforms
4
active tasks
0
idle
topic6 partitionssink tasks4 runningElasticsearch
throughput48 MB/s
Safe

All 4 tasks have partitions to read. Connect moves data without transformation, which is what makes it fast and boring — the right property for a pipeline.

A sink connector is a consumer group, so the same rule applies: tasks beyond the partition count sit idle. Source connectors are bounded by what the upstream system can be split into instead.
Connect cluster — workers with the same group.id
Worker 1REST API (8083)Worker 2REST APIWorker 3REST APIConnector Atask 0Connector Atask 1Connector Btask 0internal Kafka topicsthe workers’ stateoffset.storage.topic — source connector logical offsetsconfig.storage.topic — connector configurationsstatus.storage.topic — connector status

plugin.path is where each worker loads its connectors, converters, transformations, and secret providers from.

Source flow
converterexternal systemsource taskworkerKafka
Sink flow
converterKafkaworkersink taskexternal system
Figure 9.4.1Kafka Connect architecture

The vocabulary, precisely:

TermDefinition
Connector plug-in"libraries that Kafka Connect executes and that are responsible for moving the data"
WorkerA process in the Connect cluster. You install plug-ins on workers
ConnectorA configured instance — created/managed via REST API
Task"Connectors start additional tasks to move large amounts of data IN PARALLEL and use the available resources on the worker nodes more efficiently"
Converter"support storing those data objects in Kafka in different formats"

Converter availability: "JSON format support is part of Apache Kafka, and the Confluent Schema Registry provides Avro, Protobuf, and JSON Schema converters." → "This allows users to choose the format in which data is stored in Kafka independent of the connectors they use."

4.1 Running Connect

"Kafka Connect ships with Apache Kafka, so there is no need to install it separately. For production use, especially if you are planning to use Connect to move large amounts of data or run many connectors, YOU SHOULD RUN CONNECT ON SEPARATE SERVERS FROM YOUR KAFKA BROKERS."

(Consistent with Ch. 2: never colocate a significant application with a broker — it competes for page cache.)

bin/connect-distributed.sh config/connect-distributed.properties

Key worker configurations

ConfigNotes
bootstrap.servers"You don't need to specify every broker... but it's recommended to specify at least three."
group.id"All workers with the same group ID are part of the same Connect cluster. A connector started on the cluster will run on any worker, and so will its tasks."
plugin.pathWhere connectors, converters, transformations, and secret providers live — see below
key.converter / value.converterDefault: JSON via JSONConverter (in Apache Kafka). Or AvroConverter, ProtobufConverter, JsonSchemaConverter (Confluent Schema Registry)
rest.host.name / rest.port"Connectors are typically configured and monitored through the REST API"
⚠️ plugin.path layout — the classpath trap
plugin.path=/opt/connectors,/home/gwenshap/connectors

/opt/connectors/
├── jdbc/        ← one SUBDIRECTORY per connector
│   ├── kafka-connect-jdbc-*.jar
│   └── <all its dependencies>
└── elastic/
    ├── kafka-connect-elasticsearch-*.jar
    └── <all its dependencies>

Exception: “If the connector ships as an uberjar and has no dependencies, it can be placed directly in plugin.path and doesn't require a subdirectory.”

⚠ “But note that placing dependencies in the top-level path will not work.”

⚠ The alternative — adding connectors to the Kafka Connect classpath — is “not recommended and can introduce errors if you use a connector that brings a dependency that conflicts with one of Kafka's dependencies.”

Converter-specific configuration — the prefix rule

Prefix converter params with key.converter. or value.converter.

JSON — schema or schema-less:
key.converter.schemas.enable = true | false
value.converter.schemas.enable = true | false
Avro — needs the registry location:
key.converter.schema.registry.url
value.converter.schema.registry.url

Verifying the cluster

$ curl http://localhost:8083/
{"version":"3.0.0-SNAPSHOT","commit":"fae0784ce32a448a",
 "kafka_cluster_id":"pfkYIGZQSXm8RylvACQHdg"}

$ curl http://localhost:8083/connector-plugins
[
  {"class":"org.apache.kafka.connect.file.FileStreamSinkConnector",   "type":"sink",  "version":"3.0.0-SNAPSHOT"},
  {"class":"org.apache.kafka.connect.file.FileStreamSourceConnector", "type":"source","version":"3.0.0-SNAPSHOT"},
  {"class":"org.apache.kafka.connect.mirror.MirrorCheckpointConnector","type":"source","version":"1"},
  {"class":"org.apache.kafka.connect.mirror.MirrorHeartbeatConnector", "type":"source","version":"1"},
  {"class":"org.apache.kafka.connect.mirror.MirrorSourceConnector",    "type":"source","version":"1"}
]

Note what's in plain Apache Kafka: file source, file sink, and the three MirrorMaker 2.0 connectors — MM2 is built on Connect (Ch. 10).

STANDALONE MODE

*"similar to distributed mode — you just run bin/connect-standalone.sh... You can also pass in a connector configuration file on the command line instead of through the REST API. In this mode, all the connectors and tasks run on the one standalone worker.

It is used in cases where connectors and tasks need to run on a SPECIFIC MACHINE (e.g., the syslog connector listens on a port, so you need to know which machines it is running on)."*


On this page