skip to content

MM2 Offset Translation and Checkpoints

Keeping consumer offsets meaningful on the target cluster through checkpoints and offset translation. Interviewers ask because mirrored records do not keep their offsets, which is exactly what makes failover hard.

part ofApache Kafkaoverview, primer and where to startread it →
on this pageshow

questions

6

What is offset translation in MirrorMaker 2, and why can't you just reuse the source cluster's consumer offsets directly on the target cluster?

level: juniorimportance: must knowfreq 70%

answer

  1. offset = per-partition local counter, not global
  2. re-produced record gets a NEW offset on target
  3. raw seek -> duplicates or data loss
  4. checkpoints map source offset -> target offset
  5. conservative: replay not skip

basics

~20 s

Offset translation maps a consumer's position on the source cluster to the equivalent position on the target cluster. You can't reuse offsets directly because the same record gets a different offset on the replicated topic, so the raw number points to the wrong place.

solid answer

~40 s

When MirrorMaker 2 (MM2) replicates a topic, it re-produces each record into a new topic on the target cluster. Offsets are per-partition monotonic counters assigned at write time, so a record at offset 100 on the source may land at offset 350 on the target (the target partition may already hold older data, compacted gaps, or a different starting point). A consumer that fails over and naively seeks to source offset 100 would skip or replay messages. Offset translation solves this: the MirrorCheckpointConnector periodically records, for each consumer group, the mapping between a committed source offset and the corresponding target offset (derived from the offset-sync stream). On failover, tooling like RemoteClusterUtils.translateOffsets converts the group's source offsets into correct target offsets so consumers resume near where they left off.

go deeper

for a junior

Know that offsets differ between clusters and that translation maps source position to target position so consumers can fail over.

for a middle

Explain the re-produce mechanism, the checkpoints topic, and that translation is approximate (replay, not skip).

for a senior

Tie offset syncs -> checkpoints -> RemoteClusterUtils together and reason about translation lag and at-least-once semantics.

for a principal

Frame failover RPO/RTO tradeoffs and where translation imprecision interacts with consumer idempotency design.

## The core problem In Apache Kafka, every message in a partition gets an **offset** — a monotonically increasing integer assigned by the broker at the moment the record is appended to that partition's log. Offsets are **local to a single partition on a single cluster**; they are not global IDs and carry no meaning across clusters. A **consumer group** tracks its progress by committing the offset of the next record it wants to read, stored in the internal `__consumer_offsets` topic on that cluster. ## Why raw offsets don't transfer **MirrorMaker 2 (MM2)** is Kafka's replication tool built on Kafka Connect. To replicate topic `orders` from cluster A (source) to cluster B (target), MM2 runs a `MirrorSourceConnector` that *consumes* from `A.orders` and *re-produces* each record into `B.A.orders` (the default `DefaultReplicationPolicy` prefixes the source alias). Because the target topic is a brand-new log being written to independently, the offset a record receives on B is almost never the same number it had on A: - B's partition may have started at a different base offset. - Replication may have begun after the topic already had data. - Log compaction or retention deletion creates gaps differently on each side. - Producer retries / batching differences shift positions. So if consumer group `g` committed offset `100` on `A.orders-0`, and that exact record now lives at offset `350` on `B.A.orders-0`, blindly seeking to `100` on B would make the consumer reprocess everything from B's offset 100 to 350 (duplicates) or, in the reverse drift case, skip records (data loss). ## How MM2 solves it MM2 maintains an **offset-sync stream** (internal topic `mm2-offset-syncs.<target>.internal`) where the `MirrorSourceConnector` emits pairs of (source offset, target offset) as it produces records — effectively breadcrumbs of the source→target mapping. The **`MirrorCheckpointConnector`** consumes those syncs plus the source cluster's `__consumer_offsets`, and periodically writes **checkpoints** to a `<source>.checkpoints.internal` topic on the target. A checkpoint says, per consumer group + topic-partition: "the consumer's committed source offset X corresponds to target offset Y at this moment." On failover, **`RemoteClusterUtils.translateOffsets(...)`** (or the `MirrorClient`) reads the latest checkpoints and returns a map of translated target offsets, which an operator (or MM2 itself, if `sync.group.offsets.enabled=true`) applies to the consumer group on the target cluster. ## Edge cases / caveats - Translation is **approximate and conservative**: it resolves to the nearest sync point *at or before* the committed offset, so consumers may replay a few messages (at-least-once), never silently skip. - There is **translation lag**: checkpoints are emitted periodically (`emit.checkpoints.interval.seconds`, default 60s), so the translated offset can be slightly stale. - Translation only works for groups whose offsets are actually being checkpointed (controlled by `groups`/`groups.exclude` and offset-sync coverage).

  • Does offset translation guarantee exactly-once on failover?
    No. It's conservative and at-least-once: it resolves to the nearest sync at or before the committed offset, so a consumer may reprocess a few messages but should not skip data. Exactly-once across clusters is not provided by MM2 offset translation.
  • Which connector produces the source-to-target offset mappings that checkpoints rely on?
    The MirrorSourceConnector emits offset syncs to the mm2-offset-syncs.<target>.internal topic as it replicates records; the MirrorCheckpointConnector consumes those syncs to build checkpoints.

saying these in an interview costs you the question

  • Claiming the same record keeps the same offset on the target cluster.
  • Saying offsets are global identifiers across clusters.
  • Claiming translation is exact / guarantees exactly-once with no replay.

context

open as a page

Describe the role of the MirrorCheckpointConnector and the __checkpoints internal topic in MM2. What exactly is stored in a checkpoint record?

level: middleimportance: must knowfreq 60%

basics

~20 s

MirrorCheckpointConnector reads source consumer-group offsets plus offset syncs and writes checkpoint records to the <source>.checkpoints.internal topic on the target. Each checkpoint stores a group, topic-partition, the upstream (source) offset, and the translated downstream (target) offset.

open as a page

A consumer group on the source has been committing offsets for hours, but RemoteClusterUtils.translateOffsets returns nothing for it on the target. What would you check?

level: middleimportance: should knowfreq 40%

basics

~20 s

Check that emit.checkpoints is enabled, the group isn't excluded by groups/groups.exclude, the <source>.checkpoints.internal topic exists with data, the offset-syncs topic has syncs covering the group's partitions, and that you're querying the right target alias and renamed topics.

open as a page

Walk through how you'd fail a consumer group over to a target cluster using RemoteClusterUtils.translateOffsets, and when you'd prefer sync.group.offsets.enabled instead.

level: seniorimportance: should knowfreq 45%

basics

~20 s

Call RemoteClusterUtils.translateOffsets to get translated target offsets for the group, commit them to the group on the target, then start consumers there pointed at the replicated topics. Prefer sync.group.offsets.enabled when you want MM2 to keep the target offsets continuously synced so failover needs no extra tooling.

open as a page

How does the OffsetSyncStore decide which offsets to record, and how do emit.checkpoints.interval plus offset-sync granularity produce translation lag and offset drift?

level: seniorimportance: should knowfreq 35%

basics

~20 s

MM2 emits offset syncs sparsely (not for every record) into the offset-syncs topic, and checkpoints are emitted on an interval. Between sync points and between emissions, the translated offset is approximate, causing translation lag and a bounded replay window on failover.

open as a page

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%

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.

open as a page