Design a multi-region active-passive Kafka DR strategy. Cover replication choice, metadata/snapshot handling, failover and failback, and how you avoid split-brain and duplicate processing.
answer
- active-passive: one writer, async-replicated standby
- MM2 (translate offsets) vs Cluster Linking (offset-preserving)
- replicate offsets + configs + ACLs, IaC topic specs
- KRaft metadata rebuilt by DR's own controller
- one active cluster → no split-brain; idempotent consumers
- failback = reverse replication then controlled cutover
basics
~20 sRun a passive DR cluster in another region, replicate topics asynchronously with MirrorMaker 2 or Cluster Linking, replicate consumer offsets and ACLs/config, and keep clients pointed at the active cluster. On disaster, translate offsets, repoint clients to DR, and run only one active cluster at a time to avoid split-brain. Failback re-syncs the original direction.
solid answer
~40 sAn active-passive design has one live (active) cluster serving producers/consumers and a passive cluster in a second region kept warm by asynchronous replication. Choose MirrorMaker 2 (open source, offset-translation via checkpoints) or Cluster Linking (offset-preserving, simpler failover). Replicate not just data but also consumer-group offsets (sync.group.offsets / checkpoints), topic configs, and ACLs (MM2 has connectors/config for ACL and config sync). Metadata about brokers/partitions is rebuilt by the DR cluster's own controller (KRaft); you also snapshot/version-control topic definitions and ACLs as IaC. Failover: detect outage, translate offsets, repoint clients (bootstrap.servers/DNS), and ensure ONLY the DR cluster is active — never produce to both for the same topic, which causes split-brain/divergence. Make consumers idempotent because translation yields at-least-once replay. Failback: reverse replication, let DR drain to primary, then cut back during a controlled window.
go deeper
Recognize that DR keeps a copy of data in another region you can switch to if the primary fails.
Describe replicating data and offsets with MM2 and repointing clients on failover.
Compare MM2 vs Cluster Linking, handle offset translation, ACL/config sync, and idempotency.
Own the end-to-end topology, RPO/RTO SLAs, split-brain prevention, failback runbook, IaC metadata, and DR game days.
## Topology: active-passive One **active** cluster (region A) takes all reads/writes. A **passive** cluster (region B) sits idle for traffic but continuously receives an **asynchronous** copy of the data. No application produces to B during normal operation. This trades a small RPO (the async lag) for clean geographic isolation and a simple consistency model. ## Replication choice - **MirrorMaker 2 (MM2)** — Kafka Connect–based, open source. Connectors: `MirrorSourceConnector` (data), `MirrorCheckpointConnector` (offset checkpoints for translation), `MirrorHeartbeatConnector` (lag). Also supports **topic config sync** and **ACL sync** so the DR topics have matching partitions/configs/permissions. Offsets need **translation** (`RemoteClusterUtils` or `sync.group.offsets.enabled`). - **Cluster Linking** (Confluent) — broker-native, **offset-preserving** (mirror partitions keep identical offsets), so consumers fail over without translation. Simpler but platform-specific. Pick based on platform and how much you value offset preservation vs open-source portability. ## What 'metadata/snapshot' means here Kafka has no single 'backup the whole cluster' button; you protect metadata in layers: - **Cluster metadata** (broker/partition assignments, topic existence) lives in the **KRaft metadata log** (`@metadata` topic, controller quorum) — the DR cluster has its **own** KRaft quorum and rebuilds runtime state; you don't copy A's controller state to B. - **Topic/config/ACL definitions** — replicate via MM2 sync AND keep them as **infrastructure-as-code** (declarative topic specs) so the DR cluster can be recreated deterministically. This is the real 'snapshot/restore of metadata.' - **Consumer offsets** — replicated via checkpoints / `sync.group.offsets` so groups resume correctly. - **Compacted state topics** (e.g., changelogs) — replicate them too; their latest-per-key snapshot rebuilds downstream state on the DR side. ## Failover sequence 1. **Detect** the primary outage (health checks, monitoring). 2. **Freeze** producers to A (they're already down). 3. **Translate offsets** for consumer groups to B (skip if Cluster Linking offset-preserving). 4. **Repoint clients** — change `bootstrap.servers`/DNS/service discovery to B; stop B's mirror from A (it's the source no longer). 5. **Promote B to active.** Only B now serves writes. ## Avoiding split-brain **Split-brain** = both clusters active and accepting writes for the same logical topic, producing diverging logs you can't cleanly merge. Prevent it by enforcing **exactly one active cluster** (active-passive, not active-active). Use a single authority for 'who is active' (DNS failover, a coordination service, or manual runbook with a lock). Never let clients write to both. Because MM2 failover is **at-least-once** (offset translation lands at-or-before the true position), make consumers **idempotent** / dedupe to tolerate replay. ## Failback After region A recovers: set up replication **B→A** to backfill the writes that happened while A was down, let A catch up, then in a controlled maintenance window stop writes on B, do a final offset translation A-ward, and repoint clients back to A. Verify lag is zero before cutover. ## Edge cases / principal concerns - Topic name prefixing under MM2 DefaultReplicationPolicy (`A.orders`) vs IdentityReplicationPolicy — affects client config and active-active loop prevention. - ACL/quota/config drift between clusters causes surprises at failover — automate sync + IaC reconciliation. - DR cluster must be capacity-matched and regularly **failover-tested** (game days); an untested DR plan is a non-plan. - Exactly-once stream processing doesn't survive a cross-cluster failover cleanly — design for at-least-once across the boundary.
- Why prefer active-passive over active-active for most DR?Active-active risks split-brain and divergence (both sides writing the same logical topic, plus MM2 replication loops). Active-passive keeps a single source of truth, a simple consistency model, and clean failover at the cost of idle standby capacity.
- Is there a single backup that snapshots an entire Kafka cluster?No. You protect data via cross-cluster replication, offsets via checkpoints, and metadata (topics/configs/ACLs) via MM2 sync plus declarative IaC. The DR cluster's KRaft controller rebuilds its own runtime metadata.
- How do you do failback safely?Reverse replication (DR→primary) to backfill writes made during the outage, wait for lag to reach zero, then in a maintenance window stop writes on DR, translate offsets back, and repoint clients to the primary.
saying these in an interview costs you the question
- Proposing active-active as the default DR without addressing split-brain/loops.
- Assuming a single 'cluster snapshot/restore' button exists.
- Forgetting to replicate consumer offsets, configs, and ACLs.
- Claiming exactly-once survives a cross-cluster failover.
- Never testing the failover (untested DR).