How do you scale a Kafka Connect distributed cluster, and how does work get distributed across workers?
answer
- same group.id = one cluster
- worker > connector > task
- task = unit of parallelism
- scale = start more workers
- tasks.max caps useful workers
basics
~10 sRun 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.
solid answer
~40 sKafka Connect in distributed mode is a cluster of worker JVM processes that share a single group.id. They coordinate through Kafka's group-membership protocol (the same machinery consumer groups use), electing a leader worker that assigns connectors and tasks to all members. You scale horizontally by launching more worker processes pointed at the same group.id, bootstrap servers, and the three internal topics (config, offset, status). When a worker joins or leaves, a rebalance redistributes work so tasks are spread roughly evenly. The real unit of parallelism is the task, not the worker: a connector declares tasks.max, and the framework creates up to that many tasks. If tasks.max is smaller than the worker count, some workers sit idle, so you size tasks.max with scaling in mind.
go deeper
Know: same group.id makes a cluster, you scale by starting more workers, and tasks (not workers) are the parallelism unit.
Explain the worker/connector/task hierarchy, the three internal topics, and why tasks.max caps how many workers help a given connector.
Discuss leader election, the group-coordination protocol shared with consumer groups, and sizing tasks.max for expected worker count and source partitionability.
Reason about cluster capacity planning across many connectors, idle-worker waste, and how task granularity interacts with downstream system throughput limits.
## What Kafka Connect is Kafka Connect is a framework for streaming data between Kafka and external systems using reusable plugins called **connectors**. A *source connector* pulls data into Kafka; a *sink connector* writes Kafka data out. Connect runs in two modes: - **Standalone**: a single process, state kept in a local file. No scaling, no fault tolerance. - **Distributed**: a cluster of **worker** processes that share work and survive individual failures. This is what production uses. ## Workers, connectors, tasks Three nested concepts: - **Worker** — an OS-level JVM process running the Connect runtime. Workers are the physical units you start and stop. - **Connector** — a logical job (e.g. "copy this MySQL table to Kafka"). A connector itself does little work; it splits the job into tasks. - **Task** — the actual unit of data-copying work and the **unit of parallelism**. A connector declares `tasks.max`, and the framework asks the connector to divide its work into at most that many tasks. So the chain is: a connector produces N tasks (N ≤ `tasks.max`), and those N tasks are what get spread across workers. ## Forming a cluster Workers become one cluster when they share the same **`group.id`** (e.g. `connect-cluster`). They must also point at the same Kafka `bootstrap.servers` and the same three **internal topics**: - `config.storage.topic` — connector/task configurations - `offset.storage.topic` — source-connector progress offsets - `status.storage.topic` — connector/task status These topics are how cluster state survives a worker dying: any surviving worker can read them. ## How work is distributed Workers join a Kafka **group** using the consumer-group coordination protocol. One worker is elected **leader**. The leader computes an assignment — which connectors and which tasks each worker should run — and distributes it. This event is a **rebalance**. The goal is a roughly even spread of tasks across workers. ## Scaling out To add capacity you simply **start another worker process** with the same `group.id`. Its arrival triggers a rebalance, and the leader hands some tasks to the new worker. To scale in, stop a worker; its tasks are reassigned to the survivors. ## The tasks.max ceiling Key edge case: **a worker with no task assigned does nothing.** If a connector's `tasks.max` is 4 but you run 10 workers, at most 4 workers do that connector's work and 6 are idle (for that connector). The maximum useful parallelism for one connector equals the number of tasks it actually creates, which is bounded by `tasks.max` **and** by what the connector can split (e.g. a source connector reading a single non-partitioned table may only ever create 1 task regardless of `tasks.max`). Size `tasks.max` ≥ the number of workers you expect to use for that connector.
- If you run 6 workers but a connector has tasks.max=2, how many workers do that connector's work?At most 2. The connector creates at most 2 tasks, so only 2 workers run them; the other 4 are idle for that connector (they may still run other connectors' tasks).
- What three internal topics must all workers in a cluster share?config.storage.topic, offset.storage.topic, and status.storage.topic. They persist cluster state so any surviving worker can recover it.
saying these in an interview costs you the question
- Saying the worker (not the task) is the unit of parallelism.
- Claiming more workers always means more throughput regardless of tasks.max.
- Confusing standalone mode (no scaling) with distributed mode.
- Thinking connectors do the data copying themselves rather than splitting into tasks.