2.6 Configuring Kafka clusters
Cluster size is bound by four things:
- Scale load across servers.
- Replication guards against data loss on single failures.
- Maintenance on Kafka or the OS while staying available.
6.1 How many brokers?
Cluster size is bound by four things:
- Disk capacity
- Replica capacity per broker
- CPU capacity
- 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 recommendation | Current recommendation | |
|---|---|---|
| Replicas per broker | no more than 4,000 | no more than 14,000 |
| Replicas per cluster | no more than 200,000 | no more than 1,000,000 |
These are replicas, not partitions — partitions × RF.
| Old official recommendation | Current recommendation (well-configured env) | |
|---|---|---|
| Replicas per broker | no more than 4,000 | no more than 14,000 |
| Replicas per cluster | no more than 200,000 | no 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.