Skip to main content

Command Palette

Search for a command to run...

Multi-Leader Replication

Published
•4 min read•View as Markdown
A

Experienced full-stack developer skilled in Node.js, React.js, and Django, with expertise in building and maintaining web applications. Proficient in using Redis DB and MongoDB for back-end development, and Antd and Material UI for front-end design. Strong knowledge of C++, Python, Rust, and Julia. Well-versed in tools such as Postman, Git, Github, and Linux.

Single Leader Replication has one major drawback: There is one leader, and all writes must go through it. If we somehow failed to connect to leaders, our write operation becomes halted.
That's why multi-leader replication comes in place.
This is an extension of Single Leader Replication. This is to allow more than one node to accept write operations in DB. Replication takes place in a similar process.
Multi-Leader Replication is also known as master-master or active/active replication. In this configuration, each leader simultaneously acts as a follower to the other leaders.

Uses-Cases

  1. Multi-Datacenter Operation:- Imagine you have a database with replicas in several different data centers. That's how you can tolerate the failure of an entire data center or perhaps you will be closer to your users.

  2. Clients With Offline Operation:- Multi-Leader Replication can be used for an application that needs to continue to work while it is disconnected from the internet. In this case, every device has a local database that acts as a leader, and there is an asynchronous multi-leader replication process.

  3. Collaborative Editing:- Real-time collaborative editing applications allow several people to edit a document simultaneously. For example, Etherpad and google docs allow multiple people to concurrently edit a text document or spreadsheet.

Handling Write Conflicts

The biggest problem with multi-leader replication is that write conflicts can occur which means that conflict resolution is required.

Synchronous vs Asynchronous Replication

In a multi-leader setup, both writes are successful, and the conflict is only detected asynchronously at some later point in time.
In principle, you can make the conflict detection synchronous, by waiting for the write to be replicated to all replicas before telling the user that the write was successful. However, by doing so, you would lose the main advantage of multi-leader replication allowing each replica to write independently.

Conflict Avoidance

The simplest strategy for dealing with conflicts is to avoid them. If an application can ensure that all writes for a particular record should be done by a particular leader.
Conflicts do not occur.

However, sometimes you might want to change the designated leader for a record perhaps, because one data center has failed and you need to reroute traffic to another data center, or perhaps because a user has moved to a different location and is now closer to a different datacenter. In this situation, conflict avoidance breaks down, and you have to deal with the possibility of concurrent writing on different leaders.

Converging Toward A Consistent State

When a write operation occurs, each leader independently processes it within its local domain. However, since the system is distributed, there can be variations in the order and timing of these writes, leading to inconsistencies between the leader nodes. To address this, leaders employ a mechanism to synchronize and converge their states over time.

The convergence process involves periodic communication between the leader nodes to exchange updates and resolve any discrepancies. This ensures that the data in each leader's domain gradually align with the changes made by others. To minimize conflicts, systems may implement conflict resolution strategies based on timestamps or other prioritization mechanisms.

The convergence process continues until all leader nodes have reached a consistent state, where all write operations have been applied uniformly across the entire system. Achieving a consistent state helps maintain data integrity and enables high availability and fault tolerance.

Custom Conflict Resolution Logic

A conflict may depend on the application most multi-leader replication tools let you write conflict resolution logic using application code. That code may be executed on write or on read.

  1. On Write: As soon as the database system detects a conflict in the log of replicated changes it calls

  2. On Read: When a conflict is detected, all the conflicting writes are stored. The next time the data is read, these multiple versions of the data are returned to the application.

Multi-Leader Replication Topologies

A replication topology describes the communication paths along which writes are propagated from one node to another.

The most general topology is all-to-all, in which every leader sends its writes to every other leader. However, more restricted topologies are also used: for example, MySQL by default supports only a circular topology, in which each node receives writes from one node and forwards those writes to one other node.

More from this blog

ARPAN DAS BLOG

15 posts

Currently working as a full stack developer in GUNISMS