skip to content

Scaling and Rebalancing Across Workers

Scaling a distributed Connect cluster and how incremental cooperative rebalancing redistributes tasks as workers join and leave. Interviewers ask because older eager rebalancing stopped every connector at once.

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

questions

6

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

level: juniorimportance: must knowfreq 70%

answer

  1. same group.id = one cluster
  2. worker > connector > task
  3. task = unit of parallelism
  4. scale = start more workers
  5. tasks.max caps useful workers

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.

solid answer

~40 s

Kafka 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

for a junior

Know: same group.id makes a cluster, you scale by starting more workers, and tasks (not workers) are the parallelism unit.

for a middle

Explain the worker/connector/task hierarchy, the three internal topics, and why tasks.max caps how many workers help a given connector.

for a senior

Discuss leader election, the group-coordination protocol shared with consumer groups, and sizing tasks.max for expected worker count and source partitionability.

for a principal

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.

context

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

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 does tasks.max interact with worker count when balancing load, and what are the failure modes of misconfiguring it?

level: middleimportance: should knowfreq 50%

basics

~20 s

tasks.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.

open as a page

What does scheduled.rebalance.max.delay.ms control, and how do you tune it for worker restarts vs. genuine failures?

level: seniorimportance: should knowfreq 45%

basics

~20 s

It'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.

open as a page

How would you design a rolling restart / upgrade of a large Connect cluster to minimize rebalancing disruption?

level: principalimportance: should knowfreq 35%

basics

~20 s

Use 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.

open as a page