How does MM2 prevent infinite replication loops in an active/active topology, and how is this tied to topic naming?
answer
- prefix encodes origin
- topicSource / upstreamTopic / isInternalTopic
- skip topics originating from target
- nested prefixes peel hops
- Identity policy defeats it
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.
solid answer
~50 sIn active/active (A↔B), a naive setup would replicate A's 'orders' to B as 'A.orders', then B replicates 'A.orders' back to A — and on forever. MM2 breaks this using the source-prefix naming from DefaultReplicationPolicy plus topic filtering. The ReplicationPolicy can read a remote name and determine its origin cluster (topicSource) and whether a given cluster is the source of a topic. MM2's `MirrorSourceConnector` uses this — via the policy's loop-detection helpers and the topic-filter — to refuse replicating a topic whose origin is the target cluster, and to avoid re-replicating already-replicated topics back toward their source. The configuration knobs are the replication policy and `topic.filter.class` (default `DefaultTopicFilter`), plus the alias-based prefix. Because cycle detection depends on parsing the prefix, IdentityReplicationPolicy (no prefix) defeats it — which is exactly why identity naming is restricted to unidirectional flows.
go deeper
Know that the source prefix encodes origin and that's what stops loops.
Explain the skip rule: don't replicate a topic back to the cluster it came from.
Detail topicSource/upstreamTopic parsing, MirrorSourceConnector's filtering, and why Identity policy breaks it.
Reason about multi-hop nested prefixes, separator-vs-real-name ambiguity, and the correctness contract a custom policy must uphold to stay loop-safe.
## The loop hazard Active/active replication means both clusters replicate to each other. Consider clusters A and B, both with a flow into the other. If A's topic `orders` is copied to B and B then copies everything (including replicas) back to A, the record set ping-pongs and grows without bound. Preventing this is **cycle detection** (loop prevention). ## How naming encodes origin With `DefaultReplicationPolicy`, every replicated topic is named `<sourceAlias><separator><originalName>`, e.g. `A.orders` on B. The prefix is not decoration — it records where the topic came from. The policy exposes parsing methods such as: - `topicSource(remoteName)` → returns the immediate source alias (`A` for `A.orders`). - `upstreamTopic(remoteName)` → strips one hop (`A.orders` → `orders`). - `isInternalTopic(name)` → flags MM2/Kafka internal topics. ## The check MM2 performs `MirrorSourceConnector` enumerates candidate topics on the source and decides what to replicate. It uses the ReplicationPolicy together with the configured topic filter to exclude: 1. Topics that, per the prefix, **originated from the target cluster** — replicating those back would close a loop. 2. MM2's own internal topics and Kafka system topics. Concretely, the policy's loop-detection logic recognizes that a topic on the source already carries the target cluster's alias as its (ultimate) origin, so it is not re-shipped to that target. This is what stops `A.orders` on B from being replicated back to A as `B.A.orders`. ## Multi-hop awareness Prefixes nest across hops: A→B→C yields `B.A.orders` on C. The policy can peel hops to find the true origin, so even chained topologies don't loop a topic back to a cluster already in its path. ## Configuration surface - `replication.policy.class` — the policy that both names and parses topics. - `replication.policy.separator` — must be consistent so parsing works. - `topic.filter.class` (default `DefaultTopicFilter`) and `topics` / `topics.exclude` — what's eligible at all. - Internal-topic filtering — excludes `__consumer_offsets`, MM2 `heartbeats`/`checkpoints`/offset-sync topics, transaction state, etc. ## Why Identity policy can't do this `IdentityReplicationPolicy` keeps the original name, so a replica of `A.orders` on B is just `orders` — indistinguishable from B's native `orders`. There is no prefix to read, so origin can't be recovered and loop prevention by name fails. That is precisely why identity naming is only safe for unidirectional/active-passive flows, and active/active mandates a prefixing (or otherwise cycle-aware) policy. ## Edge cases / gotchas - If a real topic legitimately contains the separator char in its name, parsing can misidentify origin — pick separators carefully. - Custom policies must correctly implement the source/upstream/isInternal methods or they silently reintroduce loops. - Heartbeat/checkpoint internal topics have their own handling so they aren't treated as user data and looped.
- Why does using IdentityReplicationPolicy in active/active risk an infinite loop?It strips the source prefix, so a replicated topic is indistinguishable from a native one. MM2 can't recover origin from the name, loses name-based cycle detection, and replicates the topic back to its source endlessly.
- What ReplicationPolicy methods are involved in figuring out a topic's origin?topicSource (immediate source alias), upstreamTopic (strip one hop), and isInternalTopic (exclude internal/system topics) — plus the policy's loop-detection helper used by MirrorSourceConnector.
- How are multi-hop chains kept loop-free?Prefixes nest (A→B→C gives B.A.orders); the policy peels hops to find the ultimate origin, so a topic is never shipped back to any cluster already in its path.
saying these in an interview costs you the question
- Claiming MM2 uses record headers/offsets alone for loop prevention (the topic-name prefix is the primary mechanism with DefaultReplicationPolicy)
- Saying active/active works fine with IdentityReplicationPolicy
- Forgetting internal-topic filtering as part of avoiding bad re-replication