skip to content

What is replication lag in cross-cluster Kafka mirroring, and how does it relate to RPO? How would you monitor and bound it?

level: middleimportance: must knowfreq 65%

answer

  1. lag = source offset - target offset
  2. RPO ~= lag at failure instant
  3. async never RPO=0
  4. ingest > replication throughput => blowout
  5. heartbeats topic for true e2e latency

basics

~20 s

Replication lag is how far behind the target cluster's copy is from the source. If you lose the source, any unreplicated records are lost, so lag directly sets your worst-case RPO (data-loss window). You bound it by monitoring lag/throughput and provisioning enough replication capacity.

solid answer

~40 s

Replication lag is the difference between the latest offset produced on the source topic and the offset MirrorMaker 2 has copied to the target. RPO (Recovery Point Objective) is the maximum data you're willing to lose in a disaster; in async cross-cluster replication, RPO is essentially the worst-case replication lag at the moment the source fails. Lag accumulates whenever ingest rate exceeds replication throughput — caused by under-provisioned MM2 tasks, WAN bandwidth limits, throttles set too low, or a backlog after an outage. Monitor it via MM2's `record-lag`/`replication-latency-ms` JMX metrics, consumer-group lag on the MM2 source connector, and end-to-end heartbeat latency from the `heartbeats` topic. Bound it by scaling `tasks.max`/partitions, raising or tuning replication throttles, alerting on lag SLOs, and load-testing recovery from a backlog so lag drains rather than grows unbounded.

go deeper

for a junior

Know lag = how far behind the copy is, and that bigger lag means more data lost on failover.

for a middle

Tie lag to RPO, explain why async can't be zero, and name basic monitoring (consumer lag, replication-latency-ms).

for a senior

Diagnose lag-accumulation causes and use heartbeats/JMX to set and alert on a replication-latency SLO.

for a principal

Architect for an RPO budget: capacity headroom, recovery/backlog testing, topic prioritization, and the sync-vs-async tradeoff.

**What replication is here.** Cross-cluster Kafka replication (via MirrorMaker 2) is *asynchronous*: a record is acknowledged to the producer on the *source* cluster, and only later is it copied to the *target*/DR cluster. The target is always somewhat behind. **Replication lag defined.** For each topic-partition, lag = (latest source offset) - (last offset successfully written to the target). Expressed in records it's an offset gap; expressed in time it's *replication latency* — how old the most-recently-mirrored record is. MM2 is built on Kafka Connect, so under the hood the lag is the consumer-group lag of the MM2 source task reading the source topic. **RPO defined.** *RPO (Recovery Point Objective)* is the maximum amount of data, measured in time, a business can tolerate losing in a disaster. If you fail over to the target and the target was 30 seconds behind, you lose up to 30 seconds of records — so your *actual* RPO equals the *replication lag at the instant of failure*. Async replication can never give RPO=0; only synchronous schemes (e.g. a stretched cluster across AZs with `min.insync.replicas`) approach that, at a latency cost. **Why lag accumulates (RPO blowout).** Lag grows whenever *ingest rate > replication throughput*. Causes: too few MM2 tasks (`tasks.max`) for the partition count; WAN/inter-region bandwidth saturation; a replication *throttle* (`replication.factor`-unrelated bandwidth quota) set too conservatively; a target broker that is slow or under-replicated; or a backlog left after MM2 itself was down — when MM2 restarts it must catch up on everything produced during the outage, and if catch-up throughput barely exceeds steady ingest, the backlog drains very slowly or never. **Monitoring.** - MM2 JMX metrics: `replication-latency-ms` (avg/max time from source produce to target write) and `record-age-ms`; per-connector throughput. - Source-connector *consumer-group lag* on the source cluster shows offset backlog directly. - The `heartbeats` internal topic: MM2 periodically writes timestamped heartbeats; comparing source heartbeat time to when it appears on the target gives true end-to-end latency, including idle topics. - Define an SLO (e.g. p99 replication latency < 10s) and alert on breach and on sustained positive lag *derivative* (lag trending up = capacity problem, not a blip). **Bounding lag.** - Scale out: increase `tasks.max` and ensure source topics have enough partitions for parallelism. - Provision WAN/throughput headroom; size throttles above peak ingest, not at it. - Test *recovery*: deliberately stop MM2, let a backlog build, restart, and confirm lag drains to baseline within your RPO budget — this proves catch-up throughput exceeds ingest. - Consider tiered/compacted topics and excluding non-critical topics so DR-critical data gets bandwidth priority. **Edge cases.** Idle topics show stale lag metrics unless you use heartbeats. A throttle set below steady ingest guarantees ever-growing lag. After a network partition heals, simultaneous catch-up across many topics can itself saturate the link and keep lag high.

  • Why can't asynchronous cross-cluster replication give you an RPO of zero?
    Because records are acknowledged on the source before being copied to the target, there is always a window of unreplicated data. If the source dies in that window, that data is lost. Only synchronous designs (stretched clusters with min.insync.replicas across sites) approach RPO=0, trading away latency.
  • How do the heartbeats topic and offset-syncs help you measure lag accurately?
    MM2 writes timestamped heartbeat records periodically; comparing the source heartbeat timestamp to when the heartbeat appears on the target yields true end-to-end replication latency even for idle topics. Offset-syncs map source to target offsets so consumer-group lag and failover positions can be translated correctly.

saying these in an interview costs you the question

  • Claiming MM2 replication can guarantee zero data loss / RPO=0.
  • Measuring lag only on busy topics and missing idle-topic staleness (ignoring heartbeats).
  • Setting a bandwidth throttle below steady ingest and expecting lag to stay flat.

context