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?
answer
- source alias prefix: A.orders
- DefaultReplicationPolicy renames; Identity does not
- prefix = provenance = loop break
- IdentityPolicy => manual topics.exclude
- heartbeats/checkpoints internal topics filtered
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.
solid answer
~40 sMirrorMaker 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
Know that replicated topics get a prefix like 'A.orders' and that prefix stops them being copied back.
Explain DefaultReplicationPolicy renaming and how provenance via prefix breaks bidirectional loops.
Contrast Default vs Identity policy, name the policy methods, and explain when manual topic filters become mandatory.
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.