skip to content

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%

answer

  1. syncs are SPARSE, not per-record
  2. translate = nearest sync at-or-before (rounds back)
  3. log-spaced retention covers old + recent offsets
  4. emit.checkpoints.interval.seconds=60 -> staleness
  5. lag = granularity + checkpoint staleness + commit cadence -> replay

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.

solid answer

~50 s

The MirrorSourceConnector does not record a (source,target) pair for every single record — that would be huge. It emits offset syncs at an adaptive, increasing spacing per partition (historically tunably influenced, and in newer versions retaining a logarithmic spread of older syncs so both recent and far-back offsets translate well). The OffsetSyncStore on the checkpoint side keeps these points and, to translate a committed source offset, picks the nearest sync at or before it, then derives the target offset. Because syncs are sparse, translation rounds back to a sync point (granularity-induced imprecision). Separately, checkpoints themselves are emitted only every emit.checkpoints.interval.seconds (default 60), so the published mapping can be up to one interval stale. Together these produce translation lag — the gap between a consumer's true source position and the translated target position — bounding extra replay on failover but never causing skips, since rounding is conservative (at-or-before).

go deeper

for a junior

Know that translation is approximate and failover may replay a few messages.

for a middle

Explain that syncs are sparse and checkpoints are interval-based, so the translated offset is slightly behind.

for a senior

Reason about both granularity and emission-interval contributions to lag and the at-or-before rounding rule.

for a principal

Tune sync spacing vs interval vs topic load for RPO targets and mandate consumer idempotency for the replay window.

## Two independent sources of imprecision Translation accuracy is limited by **(a) how densely offset syncs are recorded** and **(b) how often checkpoints are emitted**. Both matter. ### (a) Offset-sync granularity The **MirrorSourceConnector** replicates records and writes **offset syncs** — (source offset, target offset) pairs — to `mm2-offset-syncs.<target>.internal`. Emitting one per record is wasteful, so MM2 emits them **sparsely**. Early implementations spaced syncs by a configurable amount (`offset.lag.max` influenced how far the source could drift before a new sync was forced). Later improvements (notably KIP-545-era and subsequent fixes around the `OffsetSyncStore`) retain a **spread of syncs at exponentially increasing distances into the past**, so both *recent* offsets and *older* offsets can still be translated with bounded error, rather than only the newest region being accurate. When the **OffsetSyncStore** translates a committed **source offset S**, it finds the **latest sync whose source offset <= S** and computes the corresponding **target offset** from that sync's mapping. Because S usually falls *between* two recorded syncs, the result is rounded back to the nearest known sync — **granularity-induced imprecision**. This rounding is deliberately **conservative**: it never overshoots S, so a failed-over consumer **replays** a little rather than **skipping** data. ### (b) Checkpoint emission interval The **MirrorCheckpointConnector** publishes checkpoints only every **`emit.checkpoints.interval.seconds`** (default **60**). So even if syncs were perfectly dense, the *published* group->offset mapping that failover tooling reads can be up to one interval **stale**. A consumer that committed progress 30s ago may have a checkpoint reflecting an older position. ## Translation lag and offset drift - **Offset drift**: because target offsets are assigned independently, the numeric gap (source offset vs target offset for the same record) varies over time and per partition. - **Translation lag**: the difference between a consumer's *actual* source position and the *translated* target position it would resume from. It is bounded by **(sync granularity)+(checkpoint staleness)+(commit cadence)** and shows up as **extra replay** after failover. ## Tuning levers and tradeoffs - Lower `emit.checkpoints.interval.seconds` -> fresher checkpoints, smaller lag, but more load on the checkpoints topic and connector. - Tighter sync spacing (more frequent syncs) -> finer translation granularity, larger offset-syncs topic. - `sync.group.offsets.interval.seconds` governs how often translated offsets are committed to the target group (when enabled), adding its own staleness term. ## Edge cases - If a committed source offset is **older than the oldest retained sync**, translation may be coarse or unavailable for that far-back position — a reason the logarithmic-spread retention exists. - Compaction on the offset-syncs and checkpoints topics keeps them bounded; misconfiguring retention can lose needed syncs. - Translation imprecision means **downstream consumers must be idempotent** to tolerate the replay window safely.

  • Why does the OffsetSyncStore round to the nearest sync at-or-before the committed offset rather than the nearest overall?
    Rounding back (never forward) guarantees the translated target position is at or before the true position, so consumers replay rather than skip — preserving at-least-once and avoiding data loss on failover.
  • What problem does retaining a logarithmically-spaced set of older offset syncs solve?
    It lets MM2 translate both very recent and much older committed offsets with bounded error using a small number of stored syncs, instead of only being accurate near the log head.

saying these in an interview costs you the question

  • Saying MM2 records an offset sync for every replicated record.
  • Claiming translation can overshoot and cause skipped messages.
  • Believing lowering emit.checkpoints.interval has no cost (it adds connector/topic load).

context