As a platform architect, how would you design rack/AZ topology, replication, and fetch strategy for a multi-AZ Kafka cluster to balance durability, availability, and cross-AZ cost? What are the principal tradeoffs?
answer
- rack=AZ, equal brokers/AZ, RF=3, min.ISR=2, unclean=false
- replication cross-AZ is unavoidable cost of durability
- KIP-392 = the consumer-read cost win
- freshness vs cost tradeoff on follower reads
- Cruise Control for capacity-aware rebalance
basics
~20 sMap broker.rack to AZs, use RF=3 across 3 AZs with min.insync.replicas=2 and unclean.leader.election=false for durability. Enable KIP-392 fetch-from-follower (RackAwareReplicaSelector + client.rack) to cut cross-AZ read costs. The main tradeoffs are cost vs freshness and balanced AZ sizing.
solid answer
~40 sSet broker.rack to the AZ on every broker and run equal broker counts per AZ across 3 AZs, RF=3 (one replica per AZ), min.insync.replicas=2, acks=all, unclean.leader.election.enable=false — this survives a full AZ loss with no data loss and continued writes. The cost problem is cross-AZ network: replication is inherently cross-AZ (leader→2 remote followers) and you can't avoid that and keep RF=3 durability. For consumers, enable KIP-392 (replica.selector.class=RackAwareReplicaSelector + client.rack) so reads stay intra-AZ. Producers always hit the leader, so producer→leader traffic is partly cross-AZ; some teams pin producers and accept it. Tradeoffs: follower reads trade a little freshness (HW lag) for big cost savings; even AZ broker counts matter for balanced placement; RF=3 triples replication bandwidth; Cruise Control handles capacity-aware rebalancing the built-in algorithm ignores. Validate placement with kafka-topics --describe and monitor UnderReplicatedPartitions.
go deeper
Know the recipe: rack=AZ, RF=3 across 3 AZs, and that follower reads save cost.
Pair the durability config (min.ISR=2, acks=all, unclean=false) with enabling KIP-392 for consumer cost.
Articulate which traffic is optimizable (consumer reads) vs irreducible (replication) and the freshness/cost tradeoff.
Design the full topology, justify RF and rack granularity, fold in Cruise Control and cross-region strategy, and define validation/operations.
## Design goals and the tension between them Three forces pull against each other: - **Durability** — never lose committed data, even on an AZ outage. - **Availability** — keep producing/consuming through an AZ outage. - **Cost** — minimize cloud cross-AZ network charges (often the dominant Kafka bill in the cloud). ## Baseline topology 1. **broker.rack = AZ**, set on every broker (`us-east-1a/1b/1c`). Run **equal broker counts per AZ** — uneven counts produce uneven rack-aware placement and hot brokers. 2. **RF=3 across exactly 3 AZs.** Rack-aware assignment puts one replica per AZ. This is the minimum that survives one AZ loss with a surviving majority. 3. **min.insync.replicas=2 (=RF-1), acks=all, unclean.leader.election.enable=false.** Together: a write is durable once 2 of 3 replicas (in 2 AZs) have it; an AZ loss leaves ISR=2 so writes continue; only in-sync replicas can lead so no committed data is lost. This baseline is the well-known durable cloud recipe. ## The cost layer Cloud providers bill **cross-AZ data transfer**. Kafka's traffic: - **Replication (leader→follower)** is inherently cross-AZ with one-replica-per-AZ — and you *cannot* remove it without sacrificing the durability that requires replicas in distinct AZs. Accept it as the cost of durability. - **Producer→leader**: a producer not co-located with the leader pays cross-AZ on writes. Leaders are spread across AZs, so on average ~2/3 of producer traffic is cross-AZ. Options: accept it, or use producer-side strategies (rare) — generally accepted. - **Consumer reads** are the big, *optimizable* cost. Without KIP-392 every consumer reads the leader (~2/3 cross-AZ). **Enable fetch-from-follower** (`replica.selector.class=RackAwareReplicaSelector`, consumers set `client.rack`) so each consumer reads a same-AZ follower → consumer read traffic becomes ~intra-AZ. For read-heavy fanout workloads this is the single biggest cost win. ## Principal-level tradeoffs - **Freshness vs cost (KIP-392):** follower reads expose data only up to the follower's high-watermark, so they can trail by replication lag. Latency-critical consumers may opt out (leader reads); cost-sensitive bulk consumers opt in. - **RF vs bandwidth/storage:** RF=3 triples storage and replication bandwidth vs RF=1. Going RF=4+ adds little AZ-failure protection (you already survive one AZ) at real cost — usually not worth it unless tolerating 2 simultaneous failures. - **Rack granularity:** mapping rack=AZ is standard; finer rack=server-rack inside one AZ protects against rack-level (not AZ-level) failure but doesn't help AZ outages. Choose the failure domain you actually need to survive. - **Balanced sizing:** the built-in algorithm balances replicas/leaders by count, not by *capacity or load*. For heterogeneous brokers or skewed partitions, layer **Cruise Control** for capacity-aware, rack-aware rebalancing. - **Stretch beyond 3 AZs / regions:** spanning regions adds latency that can break ISR/acks=all; usually use MirrorMaker 2 for cross-region rather than stretching one cluster. ## Operational validation - `kafka-topics.sh --describe` and inspect replica placement per AZ; confirm no partition has 2 replicas in one AZ. - Alert on `UnderReplicatedPartitions`, ISR shrink, and offline partitions. - Periodically test by draining one AZ and confirming continued writes/reads. - Re-run reassignment after topology changes — rack-awareness doesn't retro-fix existing topics.
- Why can't you eliminate cross-AZ replication traffic while keeping AZ-failure durability?Surviving an AZ loss requires replicas in distinct AZs, so the leader must replicate to followers in other AZs — that traffic is cross-AZ by definition. You can optimize consumer reads (KIP-392) but replication cross-AZ cost is the irreducible price of multi-AZ durability.
- When would you choose finer rack granularity than AZ, or coarser like region?Use rack=server-rack only if intra-AZ rack failures are your real risk and you accept it doesn't help AZ outages. Avoid stretching one cluster across regions — the latency breaks ISR/acks=all; use MirrorMaker 2 for cross-region replication instead.
- What does the built-in rack-aware algorithm not handle that Cruise Control adds?The built-in algorithm balances replica and leader counts but ignores broker capacity, disk usage, and actual load. Cruise Control performs capacity- and load-aware, still rack-aware, rebalancing for heterogeneous brokers and skewed partitions.
saying these in an interview costs you the question
- Proposing to avoid all cross-AZ traffic while keeping RF=3 AZ-failure durability — replication cross-AZ is unavoidable.
- Treating KIP-392 as reducing replication/producer cost — it only optimizes consumer reads.
- Stretching a single cluster across regions for durability instead of using MirrorMaker 2.
- Setting min.insync.replicas=RF, which sacrifices availability on any single replica loss.
- Assuming the built-in algorithm is capacity/load aware — it only balances counts.