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?
answer
- syncs are SPARSE, not per-record
- translate = nearest sync at-or-before (rounds back)
- log-spaced retention covers old + recent offsets
- emit.checkpoints.interval.seconds=60 -> staleness
- lag = granularity + checkpoint staleness + commit cadence -> replay
basics
~20 sMM2 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 sThe 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
Know that translation is approximate and failover may replay a few messages.
Explain that syncs are sparse and checkpoints are interval-based, so the translated offset is slightly behind.
Reason about both granularity and emission-interval contributions to lag and the at-or-before rounding rule.
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).