Replication in Distributed Systems
Concept
In a distributed system, storing data on only one node is a catastrophic Single Point of Failure (SPOF). Replication is the process of keeping multiple, identical copies of your data across different physical machines (or even different geographic continents). It ensures High Availability, Fault Tolerance, and allows the system to scale read traffic.
The 3 Types of Replication Architectures
1. Single-Leader (Master-Slave)
- How it works: One node is designated as the Leader. All Write requests must go to the Leader. The Leader saves the data, then sends a replication stream to the Followers (Slaves). Followers can only process Read requests.
- Pros: Extremely simple to reason about. No data conflicts, because there is only one source of truth for writes.
- Cons: The Leader is a massive bottleneck. If the system receives 100,000 writes a second, the single Leader will crash. If the Leader dies, the system must pause to elect a new one.
2. Multi-Leader (Master-Master)
- How it works: There are multiple Leaders. A user can send a Write request to any Leader. That Leader saves the data locally, and then asynchronously broadcasts the write to all other Leaders.
- Pros: Solves the Write bottleneck. Excellent for Multi-Data-Center setups (e.g., US users write to the US Leader, European users write to the EU Leader).
- Cons: Conflict Resolution Nightmare. If User A updates a record on the US Leader, and User B updates the same record on the EU Leader at the exact same millisecond, the system must employ complex logic (LWW, Vector Clocks) to figure out who wins when the Leaders sync up.
3. Leaderless
- How it works: Introduced by Amazon’s Dynamo paper (used in Cassandra and DynamoDB). There are no Leaders. Every single node in the cluster is perfectly equal. A user can send a Write request to any node.
- Quorums: To ensure data safety, the system uses Quorums. If you write to Node A, Node A takes on the role of coordinator and instantly forwards the write to Node B and C. It waits for a majority of them to reply “Saved” before telling the user the write was successful.
- Pros: Ultimate High Availability. Any node can die at any time, and the system continues to accept reads and writes seamlessly.
- Cons: Prone to reading stale data. Requires constant background “Anti-Entropy” processes (like Merkle Trees) to detect and repair data drift between nodes.
Mental Model: Leaderless Architecture
Interview Questions
Q: You use a Single-Leader architecture. You configure it for Asynchronous Replication. The Leader receives a Write, saves it to its own disk, replies ‘Success’ to the user, and then immediately catches fire and dies before sending the data to the Followers. What happens?
A: You have permanently lost data. The system will automatically promote a Follower to become the new Leader to maintain availability. However, because the old Leader died before replicating, the new Leader does not have the user’s data. From the user’s perspective, the system lied to them when it said “Success”. (To fix this, you must use Synchronous Replication, which is much slower).
Q: Explain how Multi-Leader replication handles the “Split Brain” problem.
A: Split Brain occurs when the network connection between Leader A and Leader B breaks. Both Leaders think the other one died, so they both continue accepting Writes independently. When the network heals, the database is severely corrupted with conflicting data.
To prevent this, distributed systems usually rely on an odd number of nodes and a Quorum to establish a majority vote. If a network partition isolates Leader A from the rest of the cluster, Leader A realizes it does not have the majority of votes, and intentionally demotes itself to a Read-Only state or shuts down entirely to protect data integrity.