01 / 05

Why replicate Kafka data across regions?

Difficulty: 3/10
Multi-region recovery, RPO/RTO, MirrorMaker

Why Replicate Kafka Data Across Regions

Kafka data is replicated across regions for two primary reasons: disaster recovery and regional-outage tolerance. Disaster recovery means that if the primary region becomes unavailable due to a natural disaster, a power outage, a network partition, or a cloud provider incident, a secondary region has a copy of the data and can take over. Without cross-region replication, a regional outage means data loss and downtime until the region recovers. Regional-outage tolerance means that the system remains available even when an entire region fails. This is a stronger requirement than disaster recovery because it implies that the failover is fast and ideally automated. The business driver is typically a recovery point objective (RPO) and a recovery time objective (RTO): how much data can be lost, and how long can the system be down. Cross-region replication reduces both, but at a cost in latency, network bandwidth, and operational complexity. The trade-off is between the cost of replication and the cost of downtime. For a financial platform, the cost of downtime is usually far higher than the cost of replication, which is why cross-region replication is standard.

The mechanism for cross-region replication in Kafka is not Kafka's internal replication protocol. Kafka's internal replication is designed for low-latency, high-throughput replication within a single data center or across availability zones in the same region. Cross-region latency is typically tens to hundreds of milliseconds, which is too high for synchronous replication with acks=all. Instead, cross-region replication is done asynchronously with a tool like MirrorMaker 2 (MM2), which is built on Kafka Connect. MM2 consumes from the source cluster and produces to the target cluster, preserving topic names, partitions, and offsets (with some translation). The replication is asynchronous, so there is always some replication lag: the target cluster is behind the source by the amount of time it takes to replicate. This lag determines the RPO: if the primary fails, the data that was produced but not yet replicated is lost. The trade-off is between replication lag and cost. Shorter lag requires more network bandwidth and more resources; longer lag is cheaper but increases RPO. Version note: MirrorMaker 2 is the current standard; MirrorMaker 1 is deprecated. MM2 supports active-active replication, offset translation, and topic renaming. It runs on Kafka Connect, so it inherits Connect's scaling and fault tolerance. Check the version of MM2 and the Connect cluster when designing a cross-region replication topology.

A common mistake is to assume that cross-region replication gives you zero RPO. It does not, because replication is asynchronous. Another mistake is to use Kafka's internal replication across regions, which would require very high acks and would have terrible latency. A third mistake is to replicate everything without considering the cost: replicating all topics across regions is expensive, and some topics may not need it. A better approach is to replicate only the topics that are critical for DR, and to accept that non-critical topics may be unavailable during a regional failover. The trade-off is between completeness and cost. A full replication of all topics gives the strongest DR but the highest cost; a selective replication is cheaper but requires a clear understanding of which topics are critical. For a financial platform, the critical topics are usually the ones that carry transactions and balances; analytics and logs may not need cross-region replication. Version note: MM2 supports replication policies that allow you to include or exclude topics by name or regex, so you can implement selective replication. It also supports heartbeat topics to monitor replication lag, and checkpoints to translate offsets. These features are essential for a production DR setup.

javascript
  1. 1

    Cross-region replication is for disaster recovery and regional-outage tolerance.

  2. 2

    It is asynchronous; there is always replication lag, which determines RPO.

  3. 3

    Kafka's internal replication is not suitable for cross-region; use MirrorMaker 2.

  4. 4

    MM2 runs on Kafka Connect and supports active-active, offset translation, and topic selection.

  5. 5

    Replicate only critical topics to control cost; not all topics need DR.

  6. 6

    Replication lag determines RPO; shorter lag costs more network bandwidth.

  7. 7

    MM2 heartbeat topics and checkpoints are used to monitor lag and translate offsets.

Share

Share via WhatsApp, X, Facebook, LinkedIn or copy link. Open Graph preview enabled.