Multi-Cluster and Geo-Replication
Getting data between clusters and regions: MirrorMaker 2, Cluster Linking, active-passive and active-active topologies, and their trade-offs. Interviewers raise it in any multi-region or disaster-recovery design discussion.
part ofApache Kafkaoverview, primer and where to startread it →on this pageshowhide
explore
- MirrorMaker 2 Replication Flows5 questions
- MM2 Offset Translation and Checkpoints6 questions
- MM2 Heartbeats and Replication Monitoring5 questions
- Remote Topic Naming and Replication Policies6 questions
- Cluster Linking and Mirror Topics5 questions
- Active-Passive Disaster Recovery Topologies6 questions
- Active-Active Bidirectional Replication5 questions
- Stretch Clusters and Rack-Aware Placement5 questions
- Fetch From Follower and Cross-Region Reads5 questions
- Replication Topology Pitfalls and Cycles5 questions
- Multi-Region Design Tradeoffs6 questions
- Cross-Cluster Security and Connectivity6 questions
questions
65 · 12 sectionsWhat are the three core connectors that make up MirrorMaker 2, and what does each one do?
basics
~10 sMirrorMaker 2 has three connectors: MirrorSourceConnector copies topic data from the source cluster to the target; MirrorCheckpointConnector copies consumer-group offset progress; MirrorHeartbeatConnector emits periodic heartbeats to confirm the replication path is alive.
How do you control which topics and consumer groups MM2 replicates, using allow/deny filters?
basics
~10 sMM2 uses topics/topics.exclude and groups/groups.exclude regex lists (per flow). The 'topics'/'groups' allowlist says what to replicate; the 'exclude' denylist removes matches. By default internal and MM2 bookkeeping topics/groups are excluded.
What does sync.topic.configs.enabled do, and which topic configurations does MM2 keep in sync on the target?
basics
~10 sWhen sync.topic.configs.enabled=true (the default), MM2 periodically copies the source topic's configuration (like cleanup.policy, retention, etc.) onto the mirrored target topic so they stay consistent. A property filter controls which config keys are synced.
How does MirrorSourceConnector achieve parallelism, and how do tasks.max and source partitions affect replication throughput?
basics
~20 sMirrorSourceConnector divides the source topic-partitions it must replicate across tasks. Parallelism is bounded by tasks.max and by the number of source partitions — you never get more useful tasks than partitions, since each partition is handled by one task.
How would you design an active-active bidirectional MM2 topology, and how do replication flows avoid infinite loops?
basics
~20 sYou define two flows (A->B and B->A), each running its own set of MM2 connectors. Loops are avoided because mirrored topics are prefixed with the source cluster alias, and MM2's default filters exclude already-remote topics, so B never re-mirrors A's data back to A.
What is offset translation in MirrorMaker 2, and why can't you just reuse the source cluster's consumer offsets directly on the target cluster?
basics
~20 sOffset translation maps a consumer's position on the source cluster to the equivalent position on the target cluster. You can't reuse offsets directly because the same record gets a different offset on the replicated topic, so the raw number points to the wrong place.
Describe the role of the MirrorCheckpointConnector and the __checkpoints internal topic in MM2. What exactly is stored in a checkpoint record?
basics
~20 sMirrorCheckpointConnector reads source consumer-group offsets plus offset syncs and writes checkpoint records to the <source>.checkpoints.internal topic on the target. Each checkpoint stores a group, topic-partition, the upstream (source) offset, and the translated downstream (target) offset.
A consumer group on the source has been committing offsets for hours, but RemoteClusterUtils.translateOffsets returns nothing for it on the target. What would you check?
basics
~20 sCheck that emit.checkpoints is enabled, the group isn't excluded by groups/groups.exclude, the <source>.checkpoints.internal topic exists with data, the offset-syncs topic has syncs covering the group's partitions, and that you're querying the right target alias and renamed topics.
Walk through how you'd fail a consumer group over to a target cluster using RemoteClusterUtils.translateOffsets, and when you'd prefer sync.group.offsets.enabled instead.
basics
~20 sCall RemoteClusterUtils.translateOffsets to get translated target offsets for the group, commit them to the group on the target, then start consumers there pointed at the replicated topics. Prefer sync.group.offsets.enabled when you want MM2 to keep the target offsets continuously synced so failover needs no extra tooling.
How does the OffsetSyncStore decide which offsets to record, and how do emit.checkpoints.interval plus offset-sync granularity produce translation lag and offset drift?
basics
~20 sMM2 emits offset syncs sparsely (not for every record) into the offset-syncs topic, and checkpoints are emitted on an interval. Between sync points and between emissions, the translated offset is approximate, causing translation lag and a bounded replay window on failover.
MM2 Heartbeats and Replication Monitoring
all 5 MM2 Heartbeats and Replication Monitoring questions →What is the MirrorHeartbeatConnector in MirrorMaker 2, and what is the heartbeats topic used for?
basics
~10 sMirrorHeartbeatConnector periodically writes timestamped records to a 'heartbeats' topic on the source cluster. MM2 replicates them to the target, proving the replication path is alive and letting you measure how fast it flows.
Which MM2 JMX metrics would you use to monitor replication lag and latency, and what does each one mean?
basics
~10 sMM2's MirrorSourceConnector exposes per-topic-partition JMX metrics: replication-latency-ms (time from source append to target append), record-age-ms (age of a record when consumed from source), and byte/record rates. You scrape these to track lag.
Explain emit.heartbeats.interval.seconds and the related heartbeat configs. How would you tune them and what are the trade-offs?
basics
~10 semit.heartbeats.interval.seconds sets how often MirrorHeartbeatConnector writes a heartbeat record (default 1s). emit.heartbeats.enabled (default true) turns the feature on/off. Shorter interval = finer latency resolution but more overhead.
How would you use MM2 heartbeats to observe end-to-end replication latency and estimate RPO for a DR setup?
basics
~20 sConsume the replicated <source>.heartbeats topic on the target, read each record's emit timestamp, and compute now - timestamp. That delta is your end-to-end replication latency, which approximates worst-case RPO — how much data you'd lose on failover.
Heartbeats have stopped arriving on the target for a flow, but the MM2 Connect cluster looks 'up'. How do you diagnose this, and what does heartbeat staleness tell you versus rising replication-latency-ms?
basics
~20 sStale heartbeats mean the flow is broken or stalled, not merely slow. Check the heartbeat connector/task state, the source heartbeats topic, and the MirrorSourceConnector replicating it. Rising replication-latency-ms instead means the flow works but is lagging.
Remote Topic Naming and Replication Policies
all 6 Remote Topic Naming and Replication Policies questions →When MirrorMaker 2 replicates a topic with the default settings, what does the replicated topic get named on the destination cluster, and why?
basics
~20 sBy default MirrorMaker 2 prefixes the replicated topic with the source cluster's alias and a dot. A topic named 'orders' from cluster 'us-west' becomes 'us-west.orders' on the destination. This makes the topic's origin visible and prevents name clashes.
What is IdentityReplicationPolicy, and what trade-off do you accept by using it instead of DefaultReplicationPolicy?
basics
~20 sIdentityReplicationPolicy replicates topics keeping their original name (no source-cluster prefix), so 'orders' stays 'orders' on the destination. The trade-off: you lose automatic cycle detection and risk name collisions, so it's only safe for unidirectional (active/passive) replication.
How does MM2 prevent infinite replication loops in an active/active topology, and how is this tied to topic naming?
basics
~20 sMM2 prevents loops by encoding a topic's origin in its name via the source prefix. Before replicating, it checks whether a topic already originated from the destination cluster (or matches the configured source-cluster prefix) and skips it, so records don't bounce back forever.
How does replication.policy.separator work, what is its default, and why must it be consistent across an MM2 deployment?
basics
~20 sreplication.policy.separator is the character DefaultReplicationPolicy puts between the source cluster alias and the topic name. It defaults to a dot ('.'), so 'orders' becomes 'us-west.orders'. Changing it changes remote topic names; it must match everywhere or MM2 can't parse origins.
Which topics does MM2 deliberately NOT replicate as ordinary data, and how does it identify them?
basics
~20 sMM2 skips internal and system topics: Kafka's own __consumer_offsets and transaction-state topics, and MM2's own heartbeats, checkpoints, and offset-sync topics. It identifies them via the ReplicationPolicy's isInternalTopic check and the topic filter, so they aren't replicated as user data.
What is Confluent Cluster Linking, and how does it differ at a high level from a Connect-based replicator?
basics
~20 sCluster Linking is a broker-built-in feature that copies topics from a source Kafka cluster to a destination cluster over a direct link, with no separate Connect cluster. The destination gets read-only mirror topics that keep the same message offsets.
Describe the lifecycle/states of a mirror topic and what 'promote' versus 'failover' do.
basics
~20 sA mirror topic starts in an ACTIVE (mirroring, read-only) state. To make it writable you stop mirroring and convert it: 'promote' is a graceful, synced cutover (no data loss); 'failover' is an immediate cutover used when the source is unreachable, accepting possible loss of un-replicated records.
How does Cluster Linking keep consumer group offsets and ACLs in sync across clusters, and why does it matter for failover?
basics
~20 sThe link can be configured to periodically copy committed consumer-group offsets and ACLs from the source to the destination. Because offsets are byte-for-byte identical across clusters, the copied commit positions are directly valid, so consumers can resume at the same place after failover.
Walk through creating a cluster link and the key link-level configurations you'd set, including security to the source.
basics
~20 sOn the destination, you create a named link pointing at the source's bootstrap servers and security settings, then create mirror topics on that link. Key configs include the source bootstrap/security (SASL/SSL), whether to auto-create mirror topics, and the offset/ACL sync flags.
You're designing active/passive disaster recovery across two regions with Cluster Linking. Walk through the topology, failover, fail-back, and the key constraints/trade-offs you must account for.
basics
~20 sRun primary in region A, with a link in region B mirroring A's topics (read-only) plus offset and ACL sync. On a disaster, promote/failover the mirror topics in B and repoint producers/consumers there. Fail-back later means setting up a reverse link from B to A, catching it up, then cutting back.
Active-Passive Disaster Recovery Topologies
all 6 Active-Passive Disaster Recovery Topologies questions →What is an active-passive (primary/standby) disaster recovery topology in Kafka, and how does it differ from active-active?
basics
~20 sIn active-passive DR, one Kafka cluster (primary) serves all traffic while a second cluster (standby) only receives a one-way replicated copy. Clients use the standby only after failover. Active-active runs producers/consumers on both clusters at once.
How do RPO and RTO targets shape an active-passive Kafka DR design, and what drives each?
basics
~20 sRPO (Recovery Point Objective) is the max acceptable data loss, driven by replication lag — async replication means you can lose whatever hasn't reached the standby. RTO (Recovery Time Objective) is the max acceptable downtime, driven by how fast you detect, decide, and cut clients over.
On failover to the standby cluster, why are source offsets not valid, and how does MirrorMaker 2 offset translation let consumers resume correctly?
basics
~20 sA record's offset on the standby differs from its offset on the primary because replication starts at different points and may compact/skip. MirrorMaker 2's MirrorCheckpointConnector records the primary→standby offset mapping and writes translated consumer-group checkpoints so failed-over consumers resume near where they left off.
Contrast planned versus unplanned failover in an active-passive Kafka DR setup. What changes operationally and in achievable RPO?
basics
~20 sPlanned failover (maintenance, drills) lets you stop producers, drain replication so the standby fully catches up, then cut over with RPO near zero. Unplanned failover (primary crashes) gives no chance to drain, so you lose whatever wasn't yet replicated — RPO equals the lag at failure.
Walk through how a consumer application actually fails over to the standby cluster using translated checkpoints. What must the client do, and what are the pitfalls?
basics
~20 sThe consumer reconnects to the standby's bootstrap servers, looks up its group's translated offsets (via RemoteClusterUtils/MirrorClient or pre-synced __consumer_offsets), seeks each partition to those offsets, then resumes. Pitfalls: stale checkpoints causing reprocessing, topic-name prefixes, and running the same group active on both clusters.
In an active-active bidirectional MirrorMaker 2 setup between two clusters, how does the remote-topic naming convention prevent replication loops, and what would happen if you turned it off?
basics
~20 sMM2 prefixes replicated topics with the source cluster's alias (e.g. topic 'orders' from cluster A becomes 'A.orders' on cluster B). Because the prefix marks where data came from, MM2 won't replicate a topic back to the cluster it originated from, so records don't loop forever.
In active-active, both regions accept writes to the same logical entity. What write-conflict and ordering hazards arise across regions, and how do idempotency keys and dedup help?
basics
~20 sTwo regions can write conflicting updates to the same entity concurrently, and MM2 gives no cross-region ordering or conflict resolution — it just copies records. You handle conflicts at the application layer: attach idempotency keys so re-delivered or duplicate records are deduped, and use last-writer-wins, CRDTs, or entity-region affinity to resolve concurrent updates.
Explain the per-region local + aggregate consumption pattern in an active-active deployment. How should a globally-aware consumer read all data, and what are the aggregate-cluster trade-offs?
basics
~20 sEach region has local topics (written there) plus remote mirrored topics from other regions. A region-local consumer reads just local topics; a globally-aware consumer reads both local and remote (prefixed) topics — via a regex subscription or by consuming a separate aggregate cluster that holds all regions' data in one place.
How do you actually configure bidirectional MM2 between two clusters (connectors, flows, prefixes), and how do consumers fail over between regions?
basics
~20 sDefine both clusters and enable replication in both directions in the MM2 config (A->B and B->A). MM2 runs three connectors per flow: MirrorSourceConnector (copies records, applies the prefix), MirrorCheckpointConnector (translates consumer offsets), and MirrorHeartbeatConnector (liveness). For failover, consumers use the translated checkpoint offsets to resume on the other cluster.
As an architect, when would you choose active-active bidirectional MM2 over alternatives like a stretched cluster or active-passive, and what consistency limits must you communicate to product teams?
basics
~20 sChoose active-active MM2 when both regions must serve low-latency local writes and tolerate eventual, asynchronously-replicated cross-region data. Avoid it when you need strong global consistency or strict ordering — a stretched cluster gives synchronous consistency (at WAN-latency cost), and active-passive gives simpler failover without write conflicts. The key limit to communicate: it's eventually consistent, at-least-once, with no global ordering or conflict resolution.
Stretch Clusters and Rack-Aware Placement
all 5 Stretch Clusters and Rack-Aware Placement questions →What is the broker.rack configuration in Kafka, and why does it matter for a cluster spread across availability zones?
basics
~20 sbroker.rack labels each broker with its location (e.g. its availability zone). Kafka uses these labels to spread the copies of each partition across different zones, so losing one zone does not lose all copies of the data.
In a stretch cluster across 3 AZs, how should replication.factor, min.insync.replicas, and acks be set so the cluster survives the loss of one AZ without data loss?
basics
~20 sUse replication factor 3 (one replica per AZ), min.insync.replicas=2, and producers with acks=all. Then a write is only acknowledged after two zones have it, so losing one AZ still leaves a committed copy and writes keep working.
How does synchronous cross-AZ (or cross-DC) replication latency affect producer throughput and tail latency in a stretch cluster, and what levers tune that trade-off?
basics
~20 sWith acks=all the leader waits for in-sync replicas in other zones before acknowledging, so each write pays a cross-AZ round trip. That raises per-record latency. You hide it with batching (linger.ms, batch.size) and concurrency (in-flight requests), and by keeping zones close.
Explain follower fetching (KIP-392) in a multi-AZ Kafka cluster: what problem it solves, how to enable it, and its consistency implications.
basics
~20 sFollower fetching lets a consumer read from a nearby replica in its own AZ instead of always from the leader. This cuts cross-AZ network traffic and cost. You enable it with a broker replica.selector.class and a consumer client.rack matching its zone.
What is a 2.5-DC stretch cluster design, and how do observers / asymmetric quorums (e.g. ZooKeeper or KRaft witnesses) help avoid split-brain across two main datacenters?
basics
~20 sA 2.5-DC design runs Kafka across two full datacenters plus a tiny third site holding just a tie-breaker (a witness/observer). The third site holds no data but lets the cluster keep a majority quorum if one of the two main DCs fails, avoiding split-brain.
Fetch From Follower and Cross-Region Reads
all 5 Fetch From Follower and Cross-Region Reads questions →What is 'fetch from follower' in Kafka (KIP-392), and what problem does it solve in a multi-AZ or multi-region deployment?
basics
~20 sNormally Kafka consumers read only from the partition leader. Fetch-from-follower (KIP-392) lets a consumer read from a nearby replica (follower) in its own availability zone instead, cutting cross-zone network traffic and the cloud egress cost that comes with it.
Walk through the exact configuration needed to enable rack-aware fetch-from-follower, and explain what each setting does.
basics
~20 sOn every broker set broker.rack to its AZ and set replica.selector.class to RackAwareReplicaSelector. On every consumer set client.rack to its own AZ. The broker then matches client.rack to broker.rack and points the consumer at a same-AZ follower.
When a consumer reads from a follower (KIP-392), what is the upper bound on the data it can read, and what consistency/lag implications does that create?
basics
~20 sA follower can only serve records up to its own high-watermark — the committed offset it has replicated and acknowledged. If the follower lags the leader, the consumer reads slightly older data, and end-to-end latency can grow, though it never reads uncommitted records.
You run a 3-AZ Kafka cluster on AWS and your bill is dominated by cross-AZ data transfer from consumers. As a principal engineer, how do you design a fetch-from-follower rollout to actually capture the savings, and what can undermine them?
basics
~20 sTag brokers with broker.rack=AZ, enable RackAwareReplicaSelector, ensure every partition has a replica in each consumer AZ (rack-aware placement), and set each consumer's client.rack to its AZ. Savings collapse if there's no in-rack replica, racks are misconfigured, or consumers fall back to the leader.
What is the ReplicaSelector interface, and how could you implement a custom replica selector beyond RackAwareReplicaSelector?
basics
~20 sreplica.selector.class plugs in a class implementing org.apache.kafka.common.replica.ReplicaSelector. Kafka ships LeaderSelector (default) and RackAwareReplicaSelector. You can write your own select() logic — e.g., choose by latency or a custom topology — returning the replica the consumer should read from.
In an active/active Kafka replication setup with MirrorMaker 2, how does the replication policy prevent topics from being copied back and forth in an infinite loop?
basics
~10 sMirrorMaker 2 renames remote topics with the source-cluster alias as a prefix (e.g. us-east.orders). A topic that already carries another cluster's prefix is recognized as remote and is not replicated again, breaking the loop.
What is replication lag in cross-cluster Kafka mirroring, and how does it relate to RPO? How would you monitor and bound it?
basics
~20 sReplication lag is how far behind the target cluster's copy is from the source. If you lose the source, any unreplicated records are lost, so lag directly sets your worst-case RPO (data-loss window). You bound it by monitoring lag/throughput and provisioning enough replication capacity.
What is partition-count and config drift between a source and a mirrored Kafka cluster, why is it dangerous, and how does MirrorMaker 2 try to keep them in sync?
basics
~20 sDrift is when the target topic ends up with a different partition count or different topic configs than the source. It's dangerous because key-based ordering and partitioning break after failover. MM2 periodically syncs partition counts and topic configs to keep them matched.
A team runs IdentityReplicationPolicy so topic names stay identical across two clusters. What naming and double-replication hazards does this introduce, and how do you mitigate them?
basics
~20 sWithout the alias prefix, replicas share the original name, so producers can write to the 'same' topic on both clusters and replication has no way to tell origin from replica. You must use one-way flows or strict topics.exclude filters to avoid loops and collisions.
How do you throttle replication bandwidth in Kafka, and what's the tradeoff between intra-cluster replica throttling and cross-cluster MirrorMaker 2 throttling?
basics
~20 sInside one cluster you throttle replica traffic with leader.replication.throttled.rate / follower.replication.throttled.rate (set via kafka-configs / kafka-reassign-partitions). Across clusters, MM2 runs as Connect clients, so you throttle it with producer/consumer quotas and tasks.max. Throttle too low and lag/RPO grows; too high and you saturate the WAN.
In a multi-region Kafka deployment, why does setting acks=all with replicas in remote regions hurt producer latency, and what knobs trade durability for speed?
basics
~20 sacks=all waits for all in-sync replicas to confirm a write. If some replicas live in another region, every write waits for a cross-region round trip, adding latency. You can use acks=1 or fewer in-sync replicas for speed, but you risk losing data if a region fails.
What is the difference between a 'stretch' (stretched) Kafka cluster and a 'replicated' multi-cluster topology across regions, and when would you choose each?
basics
~20 sA stretch cluster is one Kafka cluster whose brokers span regions, giving one global view but paying cross-region latency on every write. A replicated topology is separate clusters per region linked by a copy tool (MirrorMaker 2). Stretch favors consistency; replicated favors latency and isolation.
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?
basics
~20 sPlace 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.
How do data-residency and sovereignty requirements (e.g., GDPR, regional data-localization laws) constrain a multi-region Kafka design?
basics
~20 sSome data legally must stay inside a specific country/region. That forbids stretching partitions or replicating those topics to other regions. You pin data to its home region, keep replicas and consumers in-region, and only move non-restricted or anonymized data across borders.
How do cross-region/cross-AZ egress charges shape Kafka multi-region architecture, and what techniques reduce that cost?
basics
~20 sCloud providers charge for data leaving a region or AZ. Kafka replication and cross-region consumers multiply traffic, so egress can dominate the bill. Reduce it with follower fetching (read locally), compression, fewer cross-region replicas, and rack-aware placement to avoid needless cross-AZ hops.
When you replicate data between two Kafka clusters with MirrorMaker 2, what does it mean to 'secure the replication link', and what are the two layers involved?
basics
~20 sMirrorMaker 2 is just a Kafka client. Securing the link means: (1) encrypt traffic in transit with TLS, and (2) authenticate MM2 to each cluster (mTLS or SASL) so only an authorized identity can read the source and write the target.
Which Kafka ACLs does MirrorMaker 2's principal need on the source and target clusters for topic + offset replication to work?
basics
~20 sOn the source: Read on the mirrored topics and the consumer group, plus Describe. On the target: Write and Create on the mirrored/internal topics, plus DescribeConfigs/AlterConfigs for config sync. MM2 also needs access to its Connect internal topics on whichever cluster hosts them.
Show how you configure MirrorMaker 2 so it uses SASL_SSL (SCRAM) to the source cluster and mTLS to the target cluster. What property prefixes make this possible?
basics
~10 sUse the cluster-alias prefixes in MM2's config: set source.security.protocol=SASL_SSL with source.sasl.* for SCRAM, and target.security.protocol=SSL with target.ssl.keystore.* for the client cert. Each cluster's connection is configured independently under its alias.
You secure a cross-cluster MM2 link with mutual TLS. How is the Kafka principal derived from the client certificate, and how do you control it with ssl.principal.mapping.rules?
basics
~20 sWith mTLS the broker takes the client certificate's full Distinguished Name (DN) as the principal by default, e.g. 'User:CN=mm2,OU=...,O=...'. ssl.principal.mapping.rules lets you rewrite that DN with regex into a shorter principal like 'User:mm2' that your ACLs reference.
When replicating between Kafka clusters in different VPCs or cloud accounts, why is the network path (PrivateLink / VPC peering) a Kafka-specific challenge, and how does advertised.listeners interact with it?
basics
~20 sKafka clients first hit a bootstrap broker, which returns the addresses in advertised.listeners; the client then reconnects directly to each broker. Over PrivateLink/peering those advertised addresses must be resolvable and routable from the remote VPC, or the bootstrap succeeds but all follow-up connections fail.