skip to content

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