skip to content

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