Learn Labs
10. Cross-Cluster Data Mirroring

10.5 Stretch clusters — the synchronous option

"Stretch clusters are fundamentally different from other multidatacenter scenarios. To start with, THEY ARE NOT MULTICLUSTER — IT IS JUST ONE CLUSTER. As a result, we don't need a mirroring process to keep two clusters in sync. Kafka's NORMAL REPLICATION mechanism is used, as usual."

DC1 / AZ-abroker(s) · rack=a · ZK nodeDC2 / AZ-bbroker(s) · rack=b · ZK nodeDC3 / AZ-cbroker(s) · rack=c · ZK nodeONE KAFKA CLUSTER, THREE DATACENTERSNormal Kafka replication (ISR) spans datacenters.
  • acks=all + min.insync.replicas=2 + rack definitions ⇒ “every write is ACKNOWLEDGED FROM AT LEAST TWO DATACENTERS”
  • ⇒ EFFECTIVELY SYNCHRONOUS CROSS-DC REPLICATION
Figure 10.5.1Stretch clusters — the synchronous option

How it achieves synchronous cross-DC durability:

"we can configure things so the acknowledgment will be sent AFTER the message is written successfully to Kafka brokers in TWO DATACENTERS. This involves using rack definitions to make sure each partition has replicas in multiple datacenters, and the use of min.insync.replicas and acks=all."

Plus follower fetching (2.4.0+): "brokers can also be configured to enable consumers to fetch from the CLOSEST replica using rack definitions. Brokers match their rack with that of the consumer to find the local replica that is most up-to-date, falling back to the leader if a suitable local replica is not available. Consumers fetching from followers in their LOCAL datacenter achieve HIGHER THROUGHPUT, LOWER LATENCY, and LOWER COST by reducing cross-datacenter traffic." (Ch. 4 §6.6, Ch. 6 §4.2.)

Advantages:

  • SYNCHRONOUS replication — “some types of business simply require that their DR site is ALWAYS 100% SYNCHRONIZED with the primary site. This is often A LEGAL REQUIREMENT and is applied to ANY DATA STORE across the company — Kafka included.” ► i.e. RPO = 0
  • “BOTH datacenters and ALL brokers in the cluster ARE USED. THERE IS NO WASTE like we saw in active-standby.”
  • Transparent client failover — no offset translation, no client restarts

Limitations:

  • “It ONLY protects from DATACENTER FAILURES, NOT any kind of APPLICATION OR KAFKA FAILURES.” — a bad config or a Kafka bug takes it ALL down
  • “This architecture DEMANDS PHYSICAL INFRASTRUCTURE THAT NOT ALL COMPANIES CAN PROVIDE.”

⚠️ Why THREE datacenters, not two — the ZooKeeper quorum argument

“ZooKeeper requires an UNEVEN NUMBER of nodes and will remain available IF A MAJORITY of the nodes are available.”

Two datacenters + an uneven node count
DC1zk1 zk2 zk3DC2zk4 zk5ALWAYS contains a MAJORITY
  • ► “if THIS datacenter is unavailable, ZOOKEEPER IS UNAVAILABLE, and KAFKA IS UNAVAILABLE.”
  • ► You have not achieved DC-failure tolerance at all.
Three datacenters
DC1zk1 zk2DC2zk3 zk4DC3zk5no single datacenter has a MAJORITY

“you can easily allocate nodes SO NO SINGLE DATACENTER HAS A MAJORITY. So, if one datacenter is unavailable, A MAJORITY OF NODES EXIST IN THE OTHER TWO, and the ZooKeeper cluster WILL REMAIN AVAILABLE. Therefore, so will the Kafka cluster.”

Figure 10.5.4⚠️ Why THREE datacenters, not two — the ZooKeeper quorum argument

Feasibility: "if you can install Kafka (and ZooKeeper) in at least three datacenters with HIGH BANDWIDTH and LOW LATENCY between them. This can be done if your company owns three buildings on the same street, or — more commonly — by using THREE AVAILABILITY ZONES inside ONE REGION of your cloud provider."

2.5 DC ARCHITECTURE

"A popular model for stretch clusters is a 2.5 DC architecture with both Kafka and ZooKeeper running in TWO datacenters, and a third '0.5' datacenter with ONE ZooKeeper node to provide quorum if a datacenter fails."

(Also possible: ZooKeeper and Kafka in two datacenters using a ZooKeeper group configuration that allows for MANUAL failover — "However, this setup is uncommon.")


On this page