skip to content

How do you actually configure bidirectional MM2 between two clusters (connectors, flows, prefixes), and how do consumers fail over between regions?

level: middleimportance: should knowfreq 45%

answer

  1. clusters=A,B; A->B.enabled + B->A.enabled
  2. 3 connectors: Source / Checkpoint / Heartbeat
  3. Source applies DefaultReplicationPolicy prefix
  4. Checkpoint -> RemoteClusterUtils.translateOffsets
  5. sync.group.offsets.enabled for auto failover

basics

~20 s

Define both clusters and enable replication in both directions in the MM2 config (A->B and B->A). MM2 runs three connectors per flow: MirrorSourceConnector (copies records, applies the prefix), MirrorCheckpointConnector (translates consumer offsets), and MirrorHeartbeatConnector (liveness). For failover, consumers use the translated checkpoint offsets to resume on the other cluster.

solid answer

~40 s

A bidirectional setup names two clusters with distinct aliases and enables both directions. In the `mm2.properties` connect-mirror config you set `clusters=A,B`, the bootstrap servers, and `A->B.enabled=true` plus `B->A.enabled=true`, with `*.topics` whitelists. Each enabled flow runs three connectors: **MirrorSourceConnector** copies records and config, applying `DefaultReplicationPolicy` so topics become `A.<topic>` / `B.<topic>`; **MirrorCheckpointConnector** writes offset checkpoints translating source consumer-group offsets to the equivalent positions on the target; **MirrorHeartbeatConnector** emits heartbeats to a `heartbeats` topic to measure liveness and lag. ACLs, topic configs, and group offsets are synced (`sync.topic.configs.enabled`, `sync.group.offsets.enabled`, `emit.checkpoints.enabled`). For **consumer failover**, an app reading on A that must move to B uses `RemoteClusterUtils.translateOffsets(...)` (or the synced `__consumer_offsets` via offset sync) to find where its group left off in the mirrored topic, so it resumes without reprocessing everything or skipping data.

go deeper

for a junior

Know you enable both A->B and B->A and that MM2 copies topics with a prefix; offsets differ between clusters so failover needs translation.

for a middle

List the three connectors and their roles, set the key config flags, and explain offset translation for failover via checkpoints / sync.group.offsets.

for a senior

Reason about checkpoint intervals and failover staleness, ACL/config sync, HA of the MM2/Connect cluster, and heartbeat-based lag monitoring.

for a principal

Design the MM2 deployment topology (dedicated vs shared Connect, HA, multi-region), governance of aliases and topic whitelists, and the failover runbook and SLOs.

## Topology of a bidirectional flow MM2 is built on Kafka Connect. You describe **clusters** and **replication flows**. A minimal `mm2.properties`: ``` clusters = A, B A.bootstrap.servers = a1:9092,a2:9092 B.bootstrap.servers = b1:9092,b2:9092 A->B.enabled = true A->B.topics = orders|payments B->A.enabled = true B->A.topics = orders|payments replication.factor = 3 sync.topic.configs.enabled = true sync.group.offsets.enabled = true emit.checkpoints.enabled = true emit.heartbeats.enabled = true ``` Both `A->B.enabled` and `B->A.enabled` true is what makes it **bidirectional / active-active**. Distinct aliases `A` and `B` are mandatory so the prefix disambiguates origin. ## The three connectors per flow Each enabled `X->Y` flow instantiates three Kafka Connect connectors: 1. **MirrorSourceConnector** — the workhorse. Consumes source topics, **re-produces** them to the target under `DefaultReplicationPolicy` (`orders` → `X.orders`), and replicates topic configurations and ACLs. This is where prefixing (and thus loop prevention) happens. 2. **MirrorCheckpointConnector** — periodically reads source consumer-group committed offsets and writes **checkpoints** to a `<source>.checkpoints.internal` topic on the target, mapping each source offset to the corresponding target offset of the mirrored topic. This is the basis for consumer failover. 3. **MirrorHeartbeatConnector** — produces records to a `heartbeats` topic at a fixed interval. Consumers/operators measure the heartbeat's end-to-end delay to monitor replication liveness and lag (`replication-latency-ms`). ## Offset translation and consumer failover The subtle problem: offset `5000` in A's `orders` is **not** the same physical offset as in B's `A.orders` — replication starts at a different base and records may be filtered, so positions diverge. To fail a consumer group over from A to B you need the **translated** offset. MM2 solves this two ways: - **Checkpoints** via MirrorCheckpointConnector, queried programmatically with `RemoteClusterUtils.translateOffsets(props, "groupId", "A", timeout)` which returns the map of target offsets to seek to. - **`sync.group.offsets.enabled=true`** automatically writes the translated offsets into the target's `__consumer_offsets` for the group, so a restarted consumer on B simply finds committed offsets and resumes. With correct translation, a failed-over consumer resumes near where it stopped — minimal reprocessing (some duplicates, since it's at-least-once) and no large gaps. ## What gets synced - **Records** (MirrorSourceConnector). - **Topic configs** (`sync.topic.configs.enabled`) — retention, partitions tracked. - **ACLs** (`sync.topic.acls.enabled`). - **Consumer-group offsets** (`sync.group.offsets.enabled` + checkpoints). ## Running MM2 You can run the dedicated `connect-mirror-maker.sh mm2.properties` driver (standalone MM2 cluster), or deploy the connectors onto an existing Kafka Connect cluster. For HA, run multiple MM2 nodes; Connect distributes connector tasks. ## Edge cases - **Same alias on both clusters** breaks prefixing and loop prevention — keep aliases unique. - **Offset translation lag**: checkpoints are periodic (`emit.checkpoints.interval.seconds`); a failover right after a burst may translate to a slightly stale offset → some reprocessing. - **New topics**: MM2 picks up topics matching the `topics` regex on its refresh interval; there's a small delay before a new topic is mirrored. - **Heartbeats topic** must exist/be allowed on both sides for monitoring to work.

  • Why can't a failed-over consumer just reuse its committed offset number from cluster A on cluster B?
    Offsets are per-topic-partition and not aligned across clusters: B's mirrored topic starts at a different base offset and may have filtered records, so position 5000 on A maps to some other position on B. MM2's checkpoints / RemoteClusterUtils.translateOffsets give the correct translated target offset to seek to.
  • What does MirrorHeartbeatConnector give you operationally?
    It emits records to a heartbeats topic at a fixed interval; by comparing produce time to arrival time on the target you measure replication liveness and end-to-end lag (replication-latency-ms), and detect a stalled flow even when no business data is flowing.

saying these in an interview costs you the question

  • Saying MM2 is a single connector rather than three (Source/Checkpoint/Heartbeat) per flow
  • Reusing raw source offsets on the target cluster after failover instead of translating them
  • Giving both clusters the same alias, breaking prefixing and loop prevention
  • Believing MM2 auto-fails-over consumers itself — it provides translated offsets; the app/operator initiates the move

context