03 / 05

Why would you run Kafka Connect in distributed mode?

Difficulty: 5/10
Connectors, CDC

Distributed Mode: Scaling, Fault Tolerance, and Coordination for Kafka Connect

You run Kafka Connect in distributed mode for the same reasons you run any distributed system: scalability and fault tolerance. In standalone mode, everything runs in a single process, and if that process dies, all connectors stop. There is no automatic rebalancing, no horizontal scaling, and no recovery of connector state beyond what the single process has on disk. Distributed mode uses the consumer group protocol to coordinate multiple workers. Connectors are divided into tasks, and tasks are assigned to workers. When a worker joins or leaves, tasks are rebalanced automatically. This means you can scale the number of workers up or down without manual intervention, and a worker failure does not cause an outage because its tasks are reassigned to surviving workers. This is the mode you use for any production deployment where reliability matters.

The mechanism that enables this is the storage of connector configurations, offsets, and task statuses in Kafka topics. In standalone mode, these are stored locally on the worker's disk. In distributed mode, they are stored in the connect-configs, connect-offsets, and connect-status topics. Because these topics are replicated and compacted, any worker can take over a task and resume from the last committed offset. This is what makes fault tolerance possible: if a worker dies, another worker reads the connector configuration and offsets from Kafka and continues where the failed worker left off. The config.storage.topic is a single-partition compacted topic because connector configurations are a single logical entity; the offset.storage.topic has many partitions because offsets are keyed by connector and source partition. This design is critical: if you auto-create these topics with default settings, you may get deletion instead of compaction, which would lose your connector configurations and offsets. Always create them manually with the correct settings.

A common mistake is to run distributed mode with a single worker. This gives you no fault tolerance benefit because if the worker dies, everything stops. You need at least two workers for fault tolerance, and more for throughput. Another mistake is to use the same group.id as a consumer group; the group.id for Connect must be unique because it is used to form the Connect cluster group. A third mistake is to ignore the config topic replication factor. If the config topic is under-replicated and a broker fails, you can lose connector configurations. The trade-off between standalone and distributed mode is operational complexity versus reliability. Standalone is simpler and fine for development, testing, or single-purpose pipelines where downtime is acceptable. Distributed is more complex to set up but is the only choice for production workloads that need to survive failures and scale horizontally. Version note: distributed mode has been the standard production mode since early Kafka Connect versions; KRaft mode in Kafka 3.x does not change the Connect worker coordination model, though it changes how Kafka metadata is managed.

javascript
  1. 1

    Distributed mode provides fault tolerance and horizontal scaling via the consumer group protocol.

  2. 2

    Connector configs, offsets, and statuses are stored in Kafka topics, not local disk.

  3. 3

    At least two workers are needed for fault tolerance; one worker is a single point of failure.

  4. 4

    config.storage.topic should be single-partition and compacted; offset and status topics can have more partitions.

  5. 5

    Create the storage topics manually with correct compaction and replication; auto-creation may use wrong settings.

  6. 6

    group.id must be unique and must not conflict with consumer group IDs.

  7. 7

    Standalone mode is simpler but has no fault tolerance; use it for development, not production.

Share

Share via WhatsApp, X, Facebook, LinkedIn or copy link. Open Graph preview enabled.