How do you decide the number of brokers in a cluster, accounting for replication factor, broker-failure headroom, and rebalance capacity?
answer
- brokers = max(capacity, RF+F, failure-headroom, partition-ceiling)
- RF >= min.insync.replicas + F
- survivors must carry full load after losing F
- target ~60-70% steady utilization
- broker.rack spreads replicas across AZs
basics
~20 sPick brokers so total disk/network/partition load fits with each broker under ~60-70% utilization, you have at least RF brokers (plus spare for min.insync.replicas), and the cluster keeps working when N brokers fail. Survivors must absorb the dead brokers' load and re-replication.
solid answer
~50 sBroker count comes from three constraints. First, capacity: total cluster disk, network, and partition-replica load divided by per-broker safe capacity (target ~60-70% so there's headroom). Second, replication and availability: you need at least RF brokers to place RF copies on distinct brokers, and to tolerate F failures while still satisfying min.insync.replicas (acks=all) you need roughly RF + F brokers and RF >= min.insync.replicas + F. Third, failure headroom: when a broker dies, its leadership and partitions shift to survivors, who must have spare CPU/network/disk to take the extra load and to re-replicate the lost copies — so steady-state utilization must leave room for (cluster_load / (brokers - F)). Also consider rack/AZ awareness (broker.rack) so a rack/AZ loss doesn't take out all replicas of a partition, and the per-broker partition ceiling. The result: round up brokers until per-broker utilization under the worst tolerated failure stays safe.
go deeper
Know you need at least RF brokers and that failures need spare capacity.
Combine capacity load with RF and a basic one-failure headroom to pick broker count.
Reason about min.insync.replicas, re-replication load on survivors, and partition ceilings together.
Design for multi-failure/AZ tolerance, set utilization targets, and validate via simulated broker/zone loss.
## The three forces setting broker count ### 1. Raw capacity Sum the cluster's demands — disk (throughput x retention x RF), network (produce + (RF-1) replication + fan-out fetch), and total partition-replicas — then divide by what one broker can *safely* carry. 'Safely' means a **target utilization** of ~60-70%, because the remaining headroom absorbs bursts, rebalances, and failure load. So `brokers >= cluster_load / (per_broker_capacity x target_utilization)`. ### 2. Replication & durability constraints - **Replication factor (RF)**: copies per partition. RF copies must sit on **distinct brokers**, so you need **>= RF brokers** at minimum (RF=3 → at least 3). - **min.insync.replicas (ISR)**: with `acks=all`, a produce only succeeds if at least `min.insync.replicas` replicas are in sync. A common safe config is RF=3, min.insync.replicas=2 — it tolerates **one** broker down and still accepts writes. To tolerate **F** failures while still writing, you need `RF >= min.insync.replicas + F` and enough brokers to place RF copies after losing F (so ~`RF + F` brokers as a floor). ### 3. Failure headroom (the load shift) When a broker fails: - Its **leaderships** move to other replicas → survivors handle more produce/fetch traffic and CPU. - Its lost **replicas** are re-created elsewhere → a burst of **re-replication** network + disk I/O. So the steady-state per-broker utilization must be low enough that after losing F brokers, the surviving `(brokers - F)` can carry the **whole** cluster load: `cluster_load / (brokers - F) <= per_broker_capacity`. This is why 60-70% steady utilization matters — losing a broker in a small cluster spikes the rest. Example: cluster needs to do work requiring 4 brokers at full tilt. To tolerate one failure with survivors at safe load, you'd run ~6 brokers so any single loss leaves 5 absorbing the load comfortably. ## Rack / availability-zone awareness Set **broker.rack** so Kafka's replica placement spreads a partition's RF copies across racks/AZs. Then a whole rack or AZ outage can't take out all replicas of any partition. This influences broker count and placement: you want enough brokers per rack/AZ that losing one zone still leaves min.insync.replicas satisfiable. With 3 AZs and RF=3, place one replica per AZ. ## Other ceilings - **Per-broker partition limit** (~1000s of replicas): if partition count is high, you may need more brokers purely to stay under the ceiling, independent of throughput. - **Operational**: more brokers = more parallelism for recovery and reassignment, smaller blast radius per failure, but more coordination overhead. ## Decision procedure 1. Compute cluster disk/network/partition demand. 2. brokers_capacity = demand / (broker_capacity x 0.65). 3. brokers_min_RF = RF (and >= RF + F for F-failure tolerance with ISR). 4. brokers_failure = demand / (broker_capacity) + F survivors headroom. 5. brokers_partition = total_replicas / per_broker_ceiling. 6. Take the **max**, align to rack/AZ count, round up. Validate by simulating a broker (or AZ) loss and checking survivors stay under safe utilization and ISR stays satisfied.
- With RF=3 and min.insync.replicas=2, how many broker failures can you tolerate and still accept acks=all writes?One. Losing one broker leaves 2 in-sync replicas, meeting min.insync.replicas=2, so writes continue. Losing a second drops ISR below 2 and acks=all produces start failing (though reads of committed data continue).
- Why does broker.rack matter for sizing, not just placement?It ensures a partition's replicas span racks/AZs so a zone outage doesn't kill all copies. To keep min.insync.replicas satisfiable through an AZ loss you need enough brokers per AZ, which raises the broker count and constrains how replicas distribute.
saying these in an interview costs you the question
- Running brokers at 90%+ utilization with no failure headroom — one loss cascades.
- Provisioning exactly RF brokers, leaving no room to re-replicate after a failure.
- Ignoring min.insync.replicas when reasoning about write availability under failure.
- Treating rack/AZ awareness as optional for HA clusters.