skip to content

Replication Topology Pitfalls and Cycles

The ways replication topologies go wrong: cycles, name collisions, double replication, lag blowout and config drift. Interviewers use it to see whether you have run replication rather than only drawn it.

part ofApache Kafkaoverview, primer and where to startread it →
on this pageshow

questions

5

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?

level: middleimportance: must knowfreq 70%

answer

  1. source alias prefix: A.orders
  2. DefaultReplicationPolicy renames; Identity does not
  3. prefix = provenance = loop break
  4. IdentityPolicy => manual topics.exclude
  5. heartbeats/checkpoints internal topics filtered

basics

~10 s

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

solid answer

~40 s

MirrorMaker 2 (MM2) uses a ReplicationPolicy, by default DefaultReplicationPolicy, which prefixes replicated topics with the source cluster alias and a separator: a topic 'orders' from cluster 'A' becomes 'A.orders' on the target. In a bidirectional A<->B flow, MM2 on each side checks whether a topic is already 'remote' to the target — meaning it originated elsewhere — via isInternalTopic / the policy's topic-naming rules and the heartbeats/checkpoints metadata. The prefix lets the reverse flow detect that 'A.orders' on B is a replica of A's topic and refuse to replicate it back to A as 'B.A.orders'. This prefix-based identity is the core cycle-prevention mechanism. IdentityReplicationPolicy (no prefix) removes this safety net, so loop prevention must then be enforced by topic-filter regexes (topics.exclude) instead.

go deeper

for a junior

Know that replicated topics get a prefix like 'A.orders' and that prefix stops them being copied back.

for a middle

Explain DefaultReplicationPolicy renaming and how provenance via prefix breaks bidirectional loops.

for a senior

Contrast Default vs Identity policy, name the policy methods, and explain when manual topic filters become mandatory.

for a principal

Reason about mesh topologies, nested prefixes, separator collisions, and the operational risk of changing the policy on a live flow.

## The problem When you replicate topics in two directions between clusters (active/active, or a mesh of more than two clusters), naive copying creates a **cycle**: A copies 'orders' to B, B sees a topic called 'orders' and copies it back to A, and records bounce forever, multiplying infinitely. ## Kafka and MirrorMaker 2 - Kafka is a distributed log; a *topic* is a named stream of records split into partitions. - *MirrorMaker 2 (MM2)* is the Kafka Connect-based tool that copies records from a *source* cluster to a *target* cluster. - Each cluster is given a short *alias* (e.g. 'us-east', 'us-west') in the MM2 config. ## ReplicationPolicy and renaming MM2 delegates topic naming to a `ReplicationPolicy`. The default, `org.apache.kafka.connect.mirror.DefaultReplicationPolicy`, *renames* every replicated topic by prepending the source alias and a separator (default '.'): topic `orders` from cluster `A` lands on cluster `B` as `A.orders`. This is called a ***remote topic***. The methods `formatRemoteTopic(sourceAlias, topic)` build the name and `topicSource(topic)` / `upstreamTopic(topic)` parse it back. ## Why the prefix breaks the loop When MM2 runs the reverse direction (B -> A), it lists B's topics and must decide which to replicate. Because `A.orders` already carries the prefix of *another* cluster, the policy recognizes it as a topic that did not originate on B. MM2 will not re-replicate a topic that originated upstream back toward its origin, so it never creates `B.A.orders` on A. **The prefix encodes provenance, and provenance is what stops the cycle.** Internal MM2 topics (`heartbeats`, `mm2-offset-syncs`, `checkpoints`) are likewise filtered. ## IdentityReplicationPolicy Some teams want the topic to keep the *same* name across clusters (e.g. for a clean migration or a single-direction flow). `IdentityReplicationPolicy` does that — no prefix. But now provenance is invisible, so the built-in loop protection is gone. You MUST then prevent cycles manually with `topics`/`topics.exclude` regex filters, or only run a single direction. Using `IdentityReplicationPolicy` in a bidirectional setup without exclusion filters is the classic way to create a **replication storm**. ## Mesh topologies With three or more clusters, the same prefix logic generalizes; in a full mesh you can get *nested* prefixes (`A.B.orders`) representing multi-hop replicas, which the policy still treats as remote and refuses to send back toward their source. ## Edge cases - A custom separator that collides with characters in real topic names can confuse parsing. - Changing the `ReplicationPolicy` or separator on a running flow re-derives different remote names and can cause duplicate topics. - Heartbeat and offset-sync internal topics must be excluded from replication or they cause their own noise.

  • What changes about cycle prevention if you switch to IdentityReplicationPolicy?
    The prefix that encoded provenance disappears, so MM2 can no longer tell a replica from an original. You must prevent loops yourself with topics/topics.exclude regex filters or by running replication in only one direction; otherwise records cycle infinitely.
  • How does MM2 handle the topics it uses internally, like heartbeats and offset syncs?
    They are treated as internal/MM2 topics and excluded from replication by the policy's isInternalTopic/isMM2InternalTopic checks, so the heartbeats, mm2-offset-syncs, and checkpoints topics are not themselves copied around the mesh.

saying these in an interview costs you the question

  • Saying MM2 prevents loops by tracking record offsets or de-duplicating message contents (it is topic-naming/provenance, not per-record dedup).
  • Claiming IdentityReplicationPolicy is just as safe for bidirectional flows (it removes the built-in loop guard).
  • Thinking the consumer offset or consumer group somehow stops the loop.

context

open as a page

What is replication lag in cross-cluster Kafka mirroring, and how does it relate to RPO? How would you monitor and bound it?

level: middleimportance: must knowfreq 65%

basics

~20 s

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

open as a page

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?

level: seniorimportance: should knowfreq 40%

basics

~20 s

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

open as a page

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?

level: seniorimportance: should knowfreq 45%

basics

~20 s

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

open as a page

How do you throttle replication bandwidth in Kafka, and what's the tradeoff between intra-cluster replica throttling and cross-cluster MirrorMaker 2 throttling?

level: seniorimportance: should knowfreq 35%

basics

~20 s

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

open as a page