skip to content

In an active-active MM2 topology, how do the DefaultReplicationPolicy topic prefixing and offset translation interact to prevent replication loops and ensure correct failover offsets in both directions?

level: principalimportance: nice to knowfreq 25%

answer

  1. two flows A<->B, each with own syncs + checkpoints
  2. DefaultReplicationPolicy prefix = origin marker = loop break
  3. IdentityReplicationPolicy = no prefix = manage loops yourself
  4. sync.group.offsets won't clobber a locally-active group
  5. fail-back needs the reverse flow to have run

basics

~20 s

DefaultReplicationPolicy prefixes replicated topics with the source alias, so a topic isn't re-replicated back to itself, breaking loops. Each direction runs its own checkpoint/offset-sync flow, so groups can be translated and failed over correctly either way.

solid answer

~40 s

In active-active, MM2 runs two replication flows (A->B and B->A). DefaultReplicationPolicy prefixes replicated topics with the originating cluster alias (e.g. A's 'orders' becomes 'A.orders' on B). Because the policy treats already-prefixed/remote topics as not-locally-originated, the reverse flow won't pick 'A.orders' on B and ship it back to A — that's how the prefix scheme prevents infinite replication loops. Each direction also has its own offset-syncs topic and its own <source>.checkpoints.internal on the receiving cluster, plus heartbeats, so RemoteClusterUtils can translate a group's offsets for whichever direction you're failing over. With sync.group.offsets.enabled in both flows, a group's translated offsets are continuously maintained on both clusters, but MM2 won't overwrite a group actively consuming locally. If you instead use IdentityReplicationPolicy (no prefix), you lose the built-in loop protection and must scope topics carefully to avoid cycles.

go deeper

for a junior

Know active-active runs two flows and replicated topics are renamed with a prefix.

for a middle

Explain how prefixing breaks loops and that each direction has its own checkpoints.

for a senior

Reason about sync.group.offsets in both flows and not clobbering locally-active groups.

for a principal

Architect fail-back correctness, replication-policy consistency, dedup strategy, and idempotency for the replay window across an active-active estate.

## Active-active topology **Active-active** means both clusters serve traffic and replicate to each other: two MM2 flows, **A->B** and **B->A**. The risks are (1) **replication loops** (a record bouncing A->B->A->... forever) and (2) **correct offset translation in either failover direction**. ## Loop prevention via DefaultReplicationPolicy The **`DefaultReplicationPolicy`** renames a replicated topic by **prefixing it with the source cluster alias and a separator** (default `.`): topic `orders` originating on cluster `A` appears on `B` as **`A.orders`**. Loop prevention works because the policy can determine a topic's **origin** from its name. The B->A flow is configured to replicate B's *local* topics; the topic `A.orders` on B is recognized as **originating from A** (its name carries A's prefix), so it is **not** selected for replication back to A. This breaks the cycle. (MM2 also supports chained prefixes like `A.B.orders` for multi-hop, and the policy's `topicSource`/`upstreamTopic` logic unwinds these.) With **`IdentityReplicationPolicy`** (replicate with the *same* name, no prefix — common when migrating or wanting transparent names), you **lose this naming-based loop protection** and must prevent cycles yourself via explicit topic allow/deny lists, because there's no prefix to mark origin. ## Offset translation in both directions Each flow is **independent and symmetric**: - A->B writes **offset syncs** to `mm2-offset-syncs.B.internal` (on B) and **checkpoints** to `A.checkpoints.internal` on **B**. - B->A writes its own syncs and `B.checkpoints.internal` on **A**. So `RemoteClusterUtils.translateOffsets(props, "A", group, ...)` reads A-direction checkpoints, and the reverse reads B-direction checkpoints. This lets you fail a group **either way**. ## sync.group.offsets in both flows Enabling **`sync.group.offsets.enabled=true`** on both flows keeps each group's translated offsets continuously committed on **both** clusters. The crucial safety property: MM2 **does not overwrite** the offsets of a group **actively consuming locally**. In active-active, a group typically runs on one side at a time per partition set; the connector defers to the live side, preventing the two flows from fighting over the same group's offsets. ## Heartbeats and observability The **MirrorHeartbeatConnector** emits to a `heartbeats` topic, replicated both ways, giving an end-to-end liveness/lag signal independent of data volume — useful to confirm both directions are healthy before relying on translation at failover. ## Design cautions - **Fail-back correctness** depends on the **reverse flow** having run and accumulated syncs/checkpoints; if you only ever ran A->B, you can't cleanly translate B->A at fail-back. - **Avoid double-processing**: consumers on both sides reading both local and replicated copies of the same logical topic can double-count; deduplicate by topic-origin (the prefix tells you). - **Replication-policy consistency**: all MM2 components and your failover tooling must agree on the policy/separator, or renamed-topic expectations diverge and translation keys won't match. - **Idempotency**: because translation carries a replay window, end-to-end exactly-once across clusters is not provided — design consumers to tolerate replays.

  • How does DefaultReplicationPolicy stop a record from looping A->B->A forever?
    It prefixes replicated topics with the source alias (orders -> A.orders on B). The B->A flow recognizes A.orders as originating from A by its name and excludes it from replicating back, breaking the cycle.
  • What do you lose by switching to IdentityReplicationPolicy, and how do you compensate?
    You lose name-based origin tracking and loop protection; you must use explicit topic allow/deny lists (and careful flow scoping) to avoid replication cycles and double-processing.

saying these in an interview costs you the question

  • Claiming a single checkpoints topic serves both directions — each direction has its own.
  • Saying IdentityReplicationPolicy keeps the automatic loop protection (it doesn't).
  • Assuming active-active gives exactly-once across clusters.

context