On failover to the standby cluster, why are source offsets not valid, and how does MirrorMaker 2 offset translation let consumers resume correctly?
answer
- Offsets are per-cluster, assigned independently
- MirrorSourceConnector writes offset-syncs topic
- MirrorCheckpointConnector translates + emits checkpoints
- RemoteClusterUtils.translateOffsets / sync.group.offsets.enabled
- Translation is approximate → at-least-once reprocessing
- Cluster Linking preserves offsets, no translation needed
basics
~20 sA record's offset on the standby differs from its offset on the primary because replication starts at different points and may compact/skip. MirrorMaker 2's MirrorCheckpointConnector records the primary→standby offset mapping and writes translated consumer-group checkpoints so failed-over consumers resume near where they left off.
solid answer
~50 sKafka offsets are per-partition positions assigned independently by each cluster. The same logical record gets offset N on the primary but some unrelated offset M on the standby, because the standby's log starts at a different base and replication may begin mid-stream. So a consumer that was at offset 1000 on the primary cannot blindly seek to 1000 on the standby. MirrorMaker 2 solves this with the MirrorCheckpointConnector: it reads the source cluster's committed consumer-group offsets, uses the offset-sync data captured by MirrorSourceConnector (stored in the internal offset-syncs topic) to translate source offsets to the corresponding target offsets, and writes Checkpoint records (and optionally directly commits __consumer_offsets) on the target. After failover, consumers use RemoteClusterUtils.translateOffsets, the MirrorClient, or MM2's automated __consumer_offsets sync to resume near their last committed position. Translation is approximate (rounds to the nearest known sync), so at-least-once consumers may reprocess a few records.
go deeper
Know that offsets differ across clusters so consumers can't blindly resume at the same number.
Explain that MM2 records a mapping and translates committed group offsets onto the standby.
Name the connectors and internal topics, explain approximate translation and at-least-once reprocessing.
Weigh MM2 translation vs Cluster Linking offset-preservation, design the cutover to bound reprocessing and avoid double-active groups.
## Why offsets don't carry across clusters In Kafka, an **offset** is a monotonically increasing integer that identifies a record's position **within one partition of one cluster**. Each cluster assigns offsets independently as records are appended to its own log. When MirrorMaker 2 copies records from the primary to the standby, the standby is just another producer target: it appends the mirrored records to its own log and assigns its **own** offsets. These will not match the source offsets, because: - The standby's mirror topic may have started replicating at a point where the source already had millions of records — the standby's offset 0 corresponds to source offset 4,000,000, say. - Producing to the target can fail and retry, or batching/timing differs, so there's no fixed arithmetic constant linking the two. So a consumer that committed **source offset 1000** cannot simply `seek(1000)` on the target — that's a different, possibly meaningless, record. ## How MM2 captures the mapping MM2 runs three connectors: - **MirrorSourceConnector** — copies the actual records. Crucially, as it produces each batch to the target, it learns the correspondence between source offset and the resulting target offset and periodically writes **offset-sync** records to an internal topic named like `mm2-offset-syncs.<target>.internal`. Each sync says "source partition P offset X landed at target offset Y." - **MirrorCheckpointConnector** — periodically reads the source cluster's committed consumer-group offsets (from the source `__consumer_offsets`), looks up the nearest offset-sync, **translates** the source committed offset into the equivalent target offset, and emits **Checkpoint** records to a `<source>.checkpoints.internal` topic on the target. With `sync.group.offsets.enabled=true` it can also write the translated offsets straight into the target's `__consumer_offsets`. - **MirrorHeartbeatConnector** — emits heartbeats to measure connectivity/lag. ## Consuming the translation at failover Three common paths: 1. **Automated**: with `sync.group.offsets.enabled=true`, the target's `__consumer_offsets` already holds translated commits, so a consumer group simply starts on the target and resumes — provided the group isn't actively committing there (only inactive/failed-over groups should be synced, controlled by `emit.checkpoints` and group filters). 2. **Programmatic**: `RemoteClusterUtils.translateOffsets(...)` or the `MirrorClient.remoteConsumerOffsets(...)` API reads checkpoints and returns the target offsets to `seek()` to before consuming. 3. **Manual tooling**: read the checkpoints topic directly. ## Why translation is approximate Offset syncs are emitted **periodically**, not per-record (controlled by `offset.lag.max` — emit a new sync when the offset has advanced more than this). Translation rounds the committed source offset down to the nearest sync point. Result: a failed-over consumer may resume **slightly before** its true last position and **reprocess a handful of records** — which is fine for **at-least-once** processing but means you must not assume exactly-once across a DR cutover. ## Edge cases and gotchas - **IdentityReplicationPolicy**: keeps topic names identical across clusters (no `primary.` prefix). Checkpoint/translation still works, and it's preferred for DR so apps don't have to rename topics on failover. - **Lag in checkpoints**: if the primary dies, the last few group commits may not have been checkpointed, so the translated position may be staler than the consumer's actual last commit — contributing to reprocessing. - **Newer alternatives**: Confluent **Cluster Linking** preserves offsets byte-for-byte (mirror topics share offsets), eliminating translation entirely — a key reason teams pick it for DR. - Don't run the same group **actively** on both clusters while syncing offsets, or you'll clobber live commits.
- Why is offset translation approximate rather than exact, and what's the consequence?Offset-sync records are emitted periodically (governed by offset.lag.max), not per record, so translation rounds to the nearest sync point. A failed-over consumer may resume slightly earlier than its true position and reprocess a few records — fine for at-least-once, but it breaks any exactly-once assumption across the cutover.
- How does Cluster Linking change the offset-translation story?Cluster Linking mirrors topics offset-preserving — the mirror topic on the destination shares the same offsets as the source. Consumer offsets can be migrated 1:1 with no translation, simplifying DR cutover compared to MM2's approximate checkpoint translation.
- Which MM2 internal topics are involved?mm2-offset-syncs.<target>.internal holds source→target offset mappings from MirrorSourceConnector; <source>.checkpoints.internal holds translated consumer-group checkpoints from MirrorCheckpointConnector; <source>.heartbeats carries heartbeats.
saying these in an interview costs you the question
- Assuming a record keeps the same offset on both clusters under MirrorMaker 2 (it does not)
- Claiming offset translation is exact / supports exactly-once across failover
- Forgetting that MM2 needs sync.group.offsets.enabled or RemoteClusterUtils to actually apply the translation
- Confusing Cluster Linking's offset-preserving mirroring with MM2's approximate translation