10.3 Multicluster architectures
Multi-cluster topology
- 6
- mirror links
Active-active — 6 mirror links. Every cluster serves local reads and writes, so users get local latency. Requires cycle prevention (MirrorMaker 2 prefixes topics with the source cluster name) and there is no conflict resolution — you must partition writes so the same key is only written in one DC.
3.1 Hub-and-spoke
Simpler variation: just two clusters — a LEADER and a FOLLOWER.
When to use it: "data is produced in multiple datacenters and some consumers need access to the entire dataset. The architecture also allows for applications in each datacenter to only process data local to that specific datacenter. But it does NOT give access to the entire dataset from every datacenter."
Benefits:
- “data is ALWAYS produced to the LOCAL datacenter”
- “events from each datacenter are only MIRRORED ONCE — to the central DC”
- single-DC apps live in that DC; multi-DC apps live in the central DC
- “Because replication ALWAYS GOES IN ONE DIRECTION and because each consumer ALWAYS READS FROM THE SAME CLUSTER, this architecture is SIMPLE TO DEPLOY, CONFIGURE, AND MONITOR.”
The drawback, illustrated with the bank example:
*"Processors in one regional datacenter CAN'T ACCESS DATA IN ANOTHER.
Suppose we are a large bank with branches in multiple cities. We store user profiles and account history in a Kafka cluster in each city, replicated to a central cluster for business analytics. When users connect to the website or visit their local branch, they are routed to their local cluster.
However, suppose that a user visits a branch in a DIFFERENT CITY. Because the user information doesn't exist in the city they are visiting, the branch will be forced to interact with a remote cluster (NOT RECOMMENDED) or HAVE NO WAY TO ACCESS THE USER'S INFORMATION (REALLY EMBARRASSING)."*
"For this reason, use of this pattern is usually limited to only parts of the dataset that can be COMPLETELY SEPARATED between regional datacenters."
Implementation note: "for each regional datacenter you need at least one mirroring process ON THE CENTRAL DATACENTER" (consistent with principle 3 — consume remotely). "If the same topic exists in multiple datacenters, you can write all the events to one topic with the same name in the central cluster, or write events from each datacenter to a separate topic."
3.2 Active-active
Benefits:
- PRIMARY: “the ability to SERVE USERS FROM A NEARBY DATACENTER, which typically has performance benefits, WITHOUT SACRIFICING FUNCTIONALITY due to limited availability of data (as we’ve seen happen in hub-and-spoke)”
- SECONDARY: “redundancy and resilience. Since EVERY datacenter has ALL the functionality, if one is unavailable you can direct users to a remaining datacenter. This type of failover only requires NETWORK REDIRECTS OF USERS, typically THE EASIEST AND MOST TRANSPARENT type of failover.”
The drawback: conflicts from asynchronous multi-master writes.
Conflict type 1 — read-your-own-write violation
"If a user sends an event to one datacenter and reads events from another, it is possible that the event they wrote HASN'T ARRIVED at the second datacenter yet. To the user, it will look like they just added a book to their wish list and clicked on the wish list, but the book isn't there."
Mitigation: "developers usually find a way to 'STICK' each user to a specific datacenter and make sure they use the same cluster most of the time (unless they connect from a remote location or the datacenter becomes unavailable)."
Conflict type 2 — genuinely conflicting writes
"An event from one datacenter says the user ordered book A, and an event from more or less the same time at a second datacenter says the same user ordered book B. After mirroring, both datacenters have both events and thus each datacenter has two CONFLICTING events."
The design questions you must answer:
- “Do we pick ONE event as the ‘correct’ one? If so, we need CONSISTENT RULES on how to pick one event SO APPLICATIONS ON BOTH DATACENTERS WILL ARRIVE AT THE SAME CONCLUSION.”
- “Do we decide that BOTH are true and simply send the user TWO BOOKS and have another department deal with returns?”
- ► AMAZON USED TO RESOLVE CONFLICTS THAT WAY, but organizations dealing with STOCK TRADES, for example, CAN’T.
"It is important to keep in mind that if you use this architecture, YOU WILL HAVE CONFLICTS AND WILL NEED TO DEAL WITH THEM."
⚠️ Avoiding infinite mirroring loops — the namespace trick
- · “each event will only be MIRRORED ONCE”
- · “each datacenter will contain BOTH SF.users AND NYC.users, which means each datacenter will have information for ALL THE USERS”
- · Consumers subscribe to
*.users“if they wish to consume all user events” - “Another way to think of this setup is to see it as A SEPARATE NAMESPACE FOR EACH DATACENTER that contains all the topics for that datacenter.” ► “Some mirroring tools like MirrorMaker PREVENT REPLICATION CYCLES USING A SIMILAR NAMING CONVENTION.”
The alternative: record headers.
*"Record headers introduced in Apache Kafka in version 0.11.0 enable events to be TAGGED WITH THEIR ORIGINATING DATACENTER. Header information may also be used to avoid endless mirroring loops and to allow processing events from different datacenters separately. You can also implement this feature by using a structured data format for the record values (Avro is our favorite example)...
⚠️ "However, this does require EXTRA EFFORT when mirroring, since NONE OF THE EXISTING MIRRORING TOOLS WILL SUPPORT YOUR SPECIFIC HEADER FORMAT."
The N² problem:
"part of the challenge of active-active mirroring, especially with more than two datacenters, is that you will need mirroring tasks for EACH PAIR of datacenters AND EACH DIRECTION. Many mirroring tools these days can share processes, for example, using the same process for all mirroring to a destination cluster."
The verdict:
"If you find ways to handle the challenges... then this architecture is HIGHLY RECOMMENDED. It is the MOST SCALABLE, RESILIENT, FLEXIBLE, AND COST-EFFECTIVE option we are aware of. So, it is well worth the effort to figure out solutions for avoiding replication cycles, keeping users mostly in the same datacenter, and handling conflicts."
3.3 Active-standby
A “COLD” copy of all applications, that admins can START UP in an emergency.
"This is often a LEGAL REQUIREMENT rather than something that the business is actually planning on doing — but you still need to be ready."
Benefits: "simplicity in setup and the fact that it can be used in pretty much any use case. You simply install a second cluster and set up a mirroring process... No need to worry about access to data, handling conflicts, and other architectural complexities."
⚠️ The two disadvantages — and the second one is brutal
- WASTE OF A GOOD CLUSTER. “a cluster that does nothing except wait around for a disaster is a waste of resources.” Attempted fixes and their problems:
- a SMALLER DR cluster → “a RISKY decision because YOU CAN’T BE SURE THAT THIS MINIMALLY SIZED CLUSTER WILL HOLD UP DURING AN EMERGENCY”
- shift READ-ONLY workloads to the DR cluster → “they are really running a small version of a HUB-AND-SPOKE architecture with a single spoke”
- FAILOVER IS MUCH HARDER THAN IT LOOKS.
“The bottom line is that IT IS CURRENTLY NOT POSSIBLE TO PERFORM CLUSTER FAILOVER IN KAFKA WITHOUT EITHER LOSING DATA OR HAVING DUPLICATE EVENTS. OFTEN BOTH. You can MINIMIZE them but NEVER FULLY ELIMINATE them.”
And the practice requirement:
"it should go without saying that whichever failover method you choose, YOUR SRE TEAM MUST PRACTICE IT ON A REGULAR BASIS. A plan that works today may stop working after an upgrade, or perhaps new use cases make the existing tooling obsolete. Once a quarter is usually the BARE MINIMUM for failover practices. Strong SRE teams practice far more frequently. Netflix's famous Chaos Monkey, a service that randomly causes disasters, is the extreme — any day may become failover practice day."