What is replication lag in cross-cluster Kafka mirroring, and how does it relate to RPO? How would you monitor and bound it?
answer
- lag = source offset - target offset
- RPO ~= lag at failure instant
- async never RPO=0
- ingest > replication throughput => blowout
- heartbeats topic for true e2e latency
basics
~20 sReplication 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 sReplication 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
Know lag = how far behind the copy is, and that bigger lag means more data lost on failover.
Tie lag to RPO, explain why async can't be zero, and name basic monitoring (consumer lag, replication-latency-ms).
Diagnose lag-accumulation causes and use heartbeats/JMX to set and alert on a replication-latency SLO.
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.