skip to content

On failover to the standby cluster, why are source offsets not valid, and how does MirrorMaker 2 offset translation let consumers resume correctly?

level: seniorimportance: must knowfreq 60%

answer

  1. Offsets are per-cluster, assigned independently
  2. MirrorSourceConnector writes offset-syncs topic
  3. MirrorCheckpointConnector translates + emits checkpoints
  4. RemoteClusterUtils.translateOffsets / sync.group.offsets.enabled
  5. Translation is approximate → at-least-once reprocessing
  6. Cluster Linking preserves offsets, no translation needed

basics

~20 s

A 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 s

Kafka 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

for a junior

Know that offsets differ across clusters so consumers can't blindly resume at the same number.

for a middle

Explain that MM2 records a mapping and translates committed group offsets onto the standby.

for a senior

Name the connectors and internal topics, explain approximate translation and at-least-once reprocessing.

for a principal

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

context