skip to content

How does Kafka's rack-aware replica assignment algorithm decide where to place replicas, and what is the relationship between replication factor and the number of racks?

level: middleimportance: should knowfreq 45%

answer

  1. rack-alternating broker list (A1,B1,C1,A2...)
  2. skip rack already used for this partition
  3. RF ≤ R → every replica distinct rack
  4. RF > R → some rack doubles up
  5. shift start per partition for leader spread

basics

~20 s

Kafka assigns each partition's replicas by cycling through racks so no two replicas share a rack until it runs out of racks. With replication factor 3 and 3 racks, each replica lands in a different rack; with fewer racks than RF, some replicas double up.

solid answer

~40 s

When all brokers have broker.rack set, Kafka uses a rack-aware assignment that interleaves brokers by rack. It picks a starting broker, then walks the broker list in an order that alternates racks, ensuring consecutive replicas of a partition go to different racks before any rack is reused. With RF ≤ number of racks, every replica of a partition sits in a distinct rack — ideal for surviving a full rack loss. With RF > number of racks, full distinctness is impossible, so it minimizes how many replicas share a rack and still spreads leadership evenly. The algorithm also balances leader counts and total replica counts across brokers, not just racks, to avoid hot brokers. Best practice in cloud: RF=3 across exactly 3 AZs so losing one AZ leaves a 2-replica majority intact.

go deeper

for a junior

Know replicas get spread one-per-rack until racks run out.

for a middle

Explain the rack-alternating list, the RF-vs-rack-count tradeoff, and why RF=R is the sweet spot.

for a senior

Connect placement to min.insync.replicas, acks=all, and quantify survival under one rack loss for RF≤R vs RF>R.

for a principal

Reason about uneven rack sizes, leader-balance vs rack-isolation tradeoffs, and where Cruise Control supplements the built-in algorithm.

## The placement problem Given N brokers grouped into R racks, and a topic with P partitions each needing RF replicas, Kafka must decide which RF brokers hold each partition — while (a) spreading replicas across racks for fault tolerance and (b) keeping load even across brokers. ## The algorithm (rack-aware) Kafka's `AdminUtils`/`ReplicationUtils` rack-aware assignment works roughly like this: 1. **Build a rack-alternating broker list.** Brokers are ordered so that walking the list cycles through racks — e.g. racks A,B,C with brokers → A1, B1, C1, A2, B2, C2. This interleaving is the core trick. 2. **Pick a starting broker** for partition 0's leader (offset by a shift each partition, so leaders spread evenly). 3. **Place RF replicas** by stepping through the rack-alternating list, skipping a candidate if its rack is already used by this partition (until you're forced to reuse a rack because RF > R). 4. **Shift the start** for the next partition so leaders and replicas rotate across brokers — preventing one broker from being leader for everything. ## RF vs rack count — the key math - **RF ≤ R (e.g. RF=3, R=3):** every replica is in a *distinct* rack. Losing one rack removes exactly one replica per partition. With RF=3 you keep 2 replicas → a majority survives, leadership fails over, and (with min.insync.replicas=2) producers using acks=all keep working. - **RF > R (e.g. RF=3, R=2):** impossible to give every replica its own rack. At least one rack holds 2 replicas of some partition. If *that* rack fails, that partition loses 2 of 3 replicas → only 1 left, below ISR-majority comfort, and below min.insync.replicas=2 → producers with acks=all stall. So RF should generally be ≥ number of racks you want to tolerate-minus-one across. - **RF = R is the sweet spot** for surviving exactly one rack failure with a surviving majority. ## Beyond racks: broker balance The algorithm doesn't only think about racks. It also spreads **leaders** evenly (so read/write load is balanced) and total **replica counts** evenly (so disk/replication load is balanced). Rack-awareness is layered on top of this balancing, not instead of it. ## What it does NOT do - It doesn't account for **broker capacity or current load** (that's what Cruise Control adds). - It doesn't fix **existing** topics — only new assignments / explicit reassignments. - It doesn't guarantee perfect balance when broker counts per rack are uneven; lopsided rack sizes produce lopsided placement.

  • Why is RF=3 across exactly 3 AZs considered the standard cloud topology?
    Each replica lands in a distinct AZ, so losing one AZ leaves 2 replicas — a majority that keeps min.insync.replicas=2 satisfied, allowing acks=all producers and consumers to continue without data loss.
  • What goes wrong if you run RF=3 across only 2 racks?
    At least one rack holds 2 of the 3 replicas for some partitions. If that rack fails, those partitions drop to 1 replica, fall below min.insync.replicas=2, and acks=all producers stall — defeating the durability goal.

saying these in an interview costs you the question

  • Saying rack-awareness ignores broker-level load balancing — it also balances leaders and replica counts.
  • Assuming RF > rack count still gives full rack isolation.
  • Believing the algorithm considers broker disk/CPU capacity (it doesn't — that's Cruise Control).

context