skip to content

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.

level: principalimportance: should knowfreq 30%

answer

  1. active-passive: one writer, async-replicated standby
  2. MM2 (translate offsets) vs Cluster Linking (offset-preserving)
  3. replicate offsets + configs + ACLs, IaC topic specs
  4. KRaft metadata rebuilt by DR's own controller
  5. one active cluster → no split-brain; idempotent consumers
  6. failback = reverse replication then controlled cutover

basics

~20 s

Run 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 s

An 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

for a junior

Recognize that DR keeps a copy of data in another region you can switch to if the primary fails.

for a middle

Describe replicating data and offsets with MM2 and repointing clients on failover.

for a senior

Compare MM2 vs Cluster Linking, handle offset translation, ACL/config sync, and idempotency.

for a principal

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).

context