skip to content

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%

answer

  1. 3 connectors: Source / Checkpoint / Heartbeat
  2. topic: <source>.checkpoints.internal (compacted)
  3. record: group + TP + upstream + downstream offset
  4. emit.checkpoints.interval.seconds = 60
  5. sync.group.offsets.enabled applies them to target

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.

solid answer

~40 s

The MirrorCheckpointConnector is one of MM2's three connectors (alongside MirrorSourceConnector and MirrorHeartbeatConnector). It periodically reads the source cluster's committed consumer-group offsets and combines them with the offset-sync stream (mm2-offset-syncs.<target>.internal) to compute, per group and topic-partition, the corresponding target offset. It emits these as Checkpoint records to the <source>.checkpoints.internal topic on the target cluster (with DefaultReplicationPolicy, e.g. us-west.checkpoints.internal). Each record carries: consumer group id, the source topic-partition, the upstream/source committed offset, the translated downstream/target offset, plus metadata and lag. The emission cadence is controlled by emit.checkpoints.interval.seconds (default 60). RemoteClusterUtils.translateOffsets and MirrorClient read this topic to drive failover; if sync.group.offsets.enabled=true, MM2 also applies the translated offsets directly to the target __consumer_offsets.

go deeper

for a junior

Know the checkpoints topic stores per-group source->target offset mappings on the target cluster.

for a middle

Distinguish the three connectors, name the checkpoints topic, and list what a Checkpoint record contains.

for a senior

Explain how OffsetSyncStore feeds translation and how sync.group.offsets.enabled applies checkpoints to target __consumer_offsets.

for a principal

Reason about compaction, group discovery cadence, and not clobbering live target consumers during bidirectional/active-active setups.

## MM2 connector trio MirrorMaker 2 runs three Kafka Connect connectors per replication flow: 1. **MirrorSourceConnector** — replicates topic data and emits **offset syncs** to `mm2-offset-syncs.<target>.internal`. Each sync is a (source offset, target offset) pair recorded as records are re-produced. 2. **MirrorCheckpointConnector** — the subject here; produces checkpoints for **consumer offset translation**. 3. **MirrorHeartbeatConnector** — emits heartbeats to monitor connectivity/lag. ## What the checkpoint connector does The `MirrorCheckpointConnector` periodically: - Reads the **source cluster's committed consumer-group offsets** (the upstream `__consumer_offsets`) for the configured groups (`groups`, `groups.exclude`). - Loads the **OffsetSyncStore** built from `mm2-offset-syncs.<target>.internal`. - For each group/topic-partition, finds the latest offset sync at or before the group's committed source offset and computes the **translated downstream (target) offset**. - Emits a **Checkpoint** record to the target topic **`<source-alias>.checkpoints.internal`** (e.g. `primary.checkpoints.internal`). ## Anatomy of a Checkpoint record A `Checkpoint` (see `org.apache.kafka.connect.mirror.Checkpoint`) contains: - **consumer group id** - **source topic-partition** (the upstream topic) - **upstreamOffset** — the committed offset on the source - **downstreamOffset** — the translated offset on the target - **metadata** and **timestamp** The record's key encodes the group + partition so the topic can be compacted to retain the latest checkpoint per (group, partition). ## How it's consumed - **`RemoteClusterUtils.translateOffsets(props, targetClusterAlias, consumerGroupId, timeout)`** reads the checkpoints topic and returns `Map<TopicPartition, OffsetAndMetadata>` of translated target offsets — used by failover tooling. - **`MirrorClient`** offers `remoteConsumerOffsets(...)` for the same data. - If **`sync.group.offsets.enabled=true`**, the checkpoint task *itself* commits the translated offsets into the target `__consumer_offsets` (cadence via `sync.group.offsets.interval.seconds`), so a failed-over consumer group already has correct offsets without external tooling. It will not overwrite an offset for a group that is *actively consuming* on the target (to avoid clobbering live progress). ## Key configs - `emit.checkpoints.enabled` (default true) - `emit.checkpoints.interval.seconds` (default 60) - `sync.group.offsets.enabled` (default false) - `sync.group.offsets.interval.seconds` (default 60) - `refresh.groups.interval.seconds` — how often new groups are discovered ## Edge cases - A group only gets checkpoints if it has committed offsets on the source **and** the relevant partitions have offset syncs covering those offsets. - The checkpoints topic is **compacted**, so it stores the latest mapping per group/partition rather than a full history. - Translation is **conservative** (nearest sync at-or-before), favoring replay over skipping.

  • What feeds the checkpoint connector its source-to-target offset mapping?
    The OffsetSyncStore, built from the mm2-offset-syncs.<target>.internal topic that the MirrorSourceConnector populates as it replicates records.
  • Why is the checkpoints topic log-compacted?
    Only the latest checkpoint per (consumer group, topic-partition) matters for failover, so compaction retains the newest mapping and discards stale ones, keeping the topic small and fast to read.

saying these in an interview costs you the question

  • Confusing the checkpoints topic with the offset-syncs topic (Source connector writes syncs; Checkpoint connector writes checkpoints).
  • Saying the checkpoints topic lives on the source cluster — it lives on the target.
  • Claiming checkpoints store raw message payloads rather than offset mappings.

context