skip to content

How do you design a Kafka topology so that committed data survives the total loss of a region, and what are the failure modes if you size the replicas/quorum wrong?

level: principalimportance: must knowfreq 55%

answer

  1. committed ⇒ in ≥2 regions
  2. 2.5 DC = odd quorum survives region loss
  3. wrong sizing → data loss OR write stall OR no quorum
  4. unclean.leader.election=false to avoid silent loss
  5. AZ-resilient ≠ region-resilient

basics

~20 s

Place replicas in at least three failure domains (often two regions plus a tiebreaker), require enough in-sync replicas in different regions, and use acks=all. If you size it wrong, losing a region can either stop all writes (no quorum/ISR) or, worse, lose acked data.

solid answer

~50 s

Region-loss durability means an acked record and the cluster's metadata both survive a region disappearing. Design points: use rack-awareness so each partition's replicas span regions; require `min.insync.replicas` such that surviving regions still hold a committed copy (with RF spread across regions and `acks=all`); and size the controller quorum (KRaft) across an odd number of failure domains so a region loss keeps a majority — hence the **2.5 DC** pattern (two data regions + a tiebreaker site). Failure modes from wrong sizing: (1) replicas/ISR only in the lost region → committed data gone; (2) `min.insync.replicas` can't be met by survivors → all writes blocked (availability outage) even though data is safe; (3) controller quorum split 2-2 with no tiebreaker → metadata quorum lost → cluster can't elect leaders. Conversely, allowing unclean leader election to recover availability risks promoting a stale replica and silently losing committed records.

go deeper

for a junior

Know replicas must live in more than one region and acks=all is needed for durable writes.

for a middle

Explain min.insync.replicas, rack-awareness, and that a commit should exist in 2+ regions.

for a senior

Size RF/min.insync.replicas so committed implies multi-region, and configure unclean leader election deliberately.

for a principal

Design the full 2.5-DC quorum + ISR layout, set RPO/RTO targets, run region-kill game days, and reason about the ISR at the instant of failure.

## Goal: survive losing an entire region Two things must survive: (a) **data** — any record that was acknowledged with strong durability, and (b) **metadata/control** — the ability to elect leaders and keep operating. ## Data durability mechanics - **Replication factor (RF)** sets how many copies exist; **rack-awareness** (`broker.rack` + rack-aware assignment) places them in different failure domains (regions/AZs). - **ISR** is the set of replicas caught up to the leader. `acks=all` acknowledges only when all ISR members have the record. - **`min.insync.replicas`** is the floor: if the ISR drops below it, the leader rejects writes (`NotEnoughReplicasException`). To survive region loss without losing acked data, a committed record must exist in a region other than the one that might fail — so you need RF spread across regions AND `min.insync.replicas` set so that 'committed' implies 'present in ≥2 regions.' - Example: RF=4 with 2 replicas per region in two regions, `min.insync.replicas=3`, `acks=all`. A commit needs 3 ISR, which forces at least one replica in each region to be in sync — so a committed record is always in both regions and survives either region's loss. (Sizing is subtle and must be validated.) ## Control-plane durability (KRaft / quorum) Kafka's metadata is a Raft log on the **controller quorum**. A quorum needs a **majority** of voters alive. With voters split 2-2 across two regions, losing one region drops you to 2 of 4 — not a majority — and the cluster can't make progress. The fix is an **odd number of failure domains**: the **2.5 data center** topology puts data brokers in two regions and a small third site (or cheap third region) hosting a tiebreaker voter, so 2-of-3 survives a region loss. ## Failure modes of wrong sizing 1. **Data loss**: all replicas (or all ISR) in the lost region. Acked records vanish. Root cause: RF/placement not actually spanning regions, or `min.insync.replicas` too low so commits happened with only same-region copies. 2. **Write stall (availability loss, data safe)**: survivors can't meet `min.insync.replicas`. The cluster blocks producers. Safer than data loss but still an outage — common when ISR shrinks after losing a region. 3. **Quorum loss**: controller voters split evenly with no tiebreaker → no majority → no leader elections, whole cluster frozen. 4. **Silent loss via unclean leader election**: `unclean.leader.election.enable=true` lets an out-of-sync replica become leader to restore availability, but it can be missing committed records → silent data loss. Keeping it `false` favors consistency over availability. ## Practical design checklist - RF and `min.insync.replicas` chosen so 'committed' ⇒ 'in ≥2 regions.' - Rack-aware placement across regions; verify with partition reassignment audits. - Odd controller quorum across ≥3 failure domains (2.5 DC) for metadata survival. - `unclean.leader.election.enable=false` for ledgers/financial data; weigh true only where availability beats correctness. - Test with real region-kill game days; reason about the ISR at the moment of failure, not the steady state. ## Edge cases - A stretch cluster across only **two** regions cannot tolerate a region loss for the controller quorum without a third site. - Async-replicated (MM2) topologies don't lose the source on a remote-region failure, but the remote copy lags, so the DR region may be behind by the replication lag — RPO is non-zero. - Even RF=3 across 3 AZs in **one** region does not survive a **region** loss — AZ resilience ≠ region resilience.

  • Why is unclean.leader.election.enable=false the safer default for region-loss scenarios?
    When true, Kafka may promote an out-of-sync replica to leader to keep the partition available, but that replica can be missing already-committed records, causing silent data loss. Setting it false means Kafka waits for an in-sync replica instead, preserving correctness at the cost of availability — the right tradeoff for ledgers and any data where losing acked records is unacceptable.
  • Why does a two-region stretch cluster need a third site for the controller quorum?
    The KRaft controller quorum needs a majority of voters to make progress. With voters split evenly across two regions, losing one region leaves exactly half — not a majority — so no leader elections can happen and the cluster freezes. A small third site (the '.5' in 2.5 DC) hosting a tiebreaker voter keeps a majority alive after either region is lost.

saying these in an interview costs you the question

  • Believing RF=3 across AZs in one region survives a region loss (it survives only AZ loss).
  • Setting acks=all but min.insync.replicas=1 and assuming region-loss durability.
  • Splitting the controller quorum 2-2 across two regions with no tiebreaker.
  • Enabling unclean leader election for critical data to 'improve uptime' without flagging silent data loss.
  • Treating async MM2 DR as zero-RPO (it lags by the replication lag).

context