Learn Labs
2. Installing Kafka

2.6 Configuring Kafka clusters

Cluster size is bound by four things:

  1. Scale load across servers.
  2. Replication guards against data loss on single failures.
  3. Maintenance on Kafka or the OS while staying available.

6.1 How many brokers?

Cluster size is bound by four things:

  1. Disk capacity
  2. Replica capacity per broker
  3. CPU capacity
  4. Network capacity
(1) Disk capacity

Need 10 TB retained and one broker stores 2 TB → minimum 5 brokers. Add replication factor 2 → the storage requirement grows +100% → minimum 10 brokers. “Increasing the replication factor increases storage requirements by at least 100%, depending on the factor chosen.”

Increasing the replication factor increases storage requirements by at least 100%, depending on the factor chosen. ("Replicas" = the number of different brokers a single partition is copied to.)

(2) Replica capacity per broker — the numbers you should memorize

A 10-broker cluster with 1,000,000 replicas (= 500,000 partitions × RF 2) is ~100,000 replicas per broker if evenly balanced → bottlenecks in the produce, consume, and controller queues.

Old official recommendationCurrent recommendation
Replicas per brokerno more than 4,000no more than 14,000
Replicas per clusterno more than 200,000no more than 1,000,000

These are replicas, not partitions — partitions × RF.

Old official recommendationCurrent recommendation (well-configured env)
Replicas per brokerno more than 4,000no more than 14,000
Replicas per clusterno more than 200,000no more than 1,000,000

"Advances in cluster efficiency have allowed Kafka to scale much larger." Note these are replicas, not partitions — partitions × RF. People routinely mis-plan by a factor of 3 here.

(3) CPU capacity

Usually not a major bottleneck, but it can be with an excessive amount of client connections and requests on a broker. Watch overall CPU usage relative to how many unique clients and consumer groups there are, and expand to meet those needs.

(4) Network capacity — the worked example

If the network interface on a single broker is used to 80% capacity at peak, and there are two consumers of that data, the consumers will not be able to keep up with peak traffic unless there are two brokers. If replication is being used, that is an additional consumer of the data that must be accounted for.

  • Broker NIC at 80% inbound peak
  • consumer 1 reads it → +80%
  • consumer 2 reads it → +80%
  • replication reads it → +80% — easy to forget this one

⇒ total demand ≫ 100% of one NIC. “If the network interface on a single broker is used to 80% capacity at peak, and there are two consumers of that data, the consumers will not be able to keep up with peak traffic unless there are two brokers. If replication is being used, that is an additional consumer of the data that must be accounted for.”

Also account for traffic not being consistent over the retention period (bursts at peak times). And you may scale out to more brokers to address lesser disk throughput or available system memory.

6.2 Broker configuration

Only two requirements to join a cluster (see §0): identical zookeeper.connect, unique broker.id. Other cluster-relevant params — specifically those controlling replication — are covered in later chapters.


On this page