6.3 Multi-Leader Replication
Both perform automatic merges for all the types above; different design philosophies and performance characteristics.
Replication topology
- 0 nodes
- cut off
Every leader sends its writes to every other. Most fault-tolerant — but messages can overtake each other, which is where the causality bug lives: an UPDATE can arrive before the INSERT it depends on, and timestamps do not fix it because those writes are causally related, not concurrent.
Conflict resolution
Last write wins — data is discarded. One device's removal simply overwrites the other's. LWW is only harmless when the writes are genuinely interchangeable; here it silently drops a user action.
Also: active/active or bidirectional replication. More than one node accepts writes; each leader simultaneously acts as a follower to the other leaders.
Why the book only discusses the asynchronous variant: with two leaders A and B, if writes are synchronously replicated from A to B and the network between them is interrupted, you can't write to A until the connection is restored. Synchronous multi-leader is therefore very similar to single-leader (equivalent to making B the leader and having A forward writes). The rest of the section is asynchronous multi-leader, in which any leader can process writes even when its connection to the other leaders is interrupted.
3.1 Geographically distributed operation
It rarely makes sense to use a multi-leader setup within a SINGLE region — the benefits rarely outweigh the added complexity.
Four-way comparison against single-leader in a multi-region deployment:
| Dimension | Single-leader | Multi-leader |
|---|---|---|
| Performance | Every write goes over the internet to the leader's region — significant added latency, might defeat the purpose of having multiple regions | Every write is processed in the LOCAL region and replicated asynchronously. The inter-region delay is hidden from users |
| Tolerance of regional outages | Failover can promote a follower in another region | Each region continues operating INDEPENDENTLY; replication catches up when the offline region returns |
| Tolerance of network problems | Very sensitive to the inter-region link — a client must send its request over that link and wait for the response | Can tolerate network problems better; during a temporary interruption each region's leader continues independently |
| Consistency | Can provide strong guarantees, e.g. serializable transactions | MUCH WEAKER — the biggest downside. |
The consistency limitation, stated precisely:
You can't guarantee that a bank account won't go negative or that a username is unique; it's always possible for different leaders to process writes that are INDIVIDUALLY FINE (paying out some of the money in an account, registering a particular username) but that VIOLATE THE CONSTRAINT WHEN TAKEN TOGETHER with another write on another leader.
This is simply a fundamental limitation of distributed systems. If you need to enforce such constraints, you're better off with a single-leader system.
Support: MySQL, Oracle, SQL Server, YugabyteDB; as an external add-on in Redis Enterprise, EDB Postgres Distributed, pglogical.
⚠️ As multi-leader replication is a RETROFITTED feature in many databases, there are often subtle configuration pitfalls and surprising interactions with other database features. Autoincrementing keys, triggers, and integrity constraints can be problematic. For this reason, multi-leader replication is often considered DANGEROUS TERRITORY that should be avoided if possible.
3.2 Replication topologies
(a) CIRCULAR (b) STAR (c) ALL-TO-ALL
(most general)
L1 ──▶ L2 L2 L1 ◀──▶ L2
▲ │ ▲ ▲ ╲ ╱ ▲
│ ▼ L1 ◀──R──▶ L3 │ ╲ ╱ │
L4 ◀── L3 │ │ ╳ │
▼ │ ╱ ╲ │
each node receives L4 ▼ ╱ ╲ ▼
from ONE node and L4 ◀──▶ L3
forwards to ONE one designated ROOT
other node forwards to all others every leader sends its
(generalizable to a TREE) writes to every otherLoop prevention (needed in circular and star, where a write passes through several nodes): each node has a unique identifier, and in the replication log each write is tagged with the identifiers of all nodes it has passed through. When a node receives a change tagged with its OWN identifier, it ignores it — it knows it already processed it.
Problems with each:
| Topology | Problem |
|---|---|
| Circular / star | If just ONE node fails, it can interrupt the flow of replication messages between other nodes, leaving them unable to communicate until it's fixed. The topology could be reconfigured to work around the failed node, but in most deployments that must be done MANUALLY |
| All-to-all | Better fault tolerance (messages travel along different paths, avoiding a single point of failure) — but some network links may be faster than others, so replication messages can OVERTAKE others |
The overtaking problem — a causality violation:
Simply attaching a TIMESTAMP to every write is NOT SUFFICIENT, because clocks cannot be trusted to be sufficiently in sync to correctly order these events (Ch 9).
To order these events correctly, VERSION VECTORS can be used (§4.6). However, many multi-leader replication systems don't use good techniques for ordering updates, leaving them vulnerable to this issue. It's worth carefully reading the documentation and THOROUGHLY TESTING your database to ensure it really does provide the guarantees you believe it has.
3.3 Sync engines and local-first software
The other major multi-leader use case: applications that must work while disconnected from the internet. Calendar apps on phone, laptop, and other devices — you must be able to read and write at any time, regardless of connectivity, and changes sync when next online.
Every device has a local database replica that acts as a LEADER (it accepts writes), and there is an asynchronous multi-leader replication process (SYNC) between all your devices. The replication lag may be HOURS OR EVEN DAYS.
Architecturally this is multi-leader replication between regions, taken to the extreme: each DEVICE is a "region," and the network connection between them is extremely unreliable.
Real-time collaboration is the same architecture. Google Docs/Sheets, Figma, Linear. What makes them responsive is that user input is immediately reflected in the UI without waiting for a network round-trip, and edits are shown to collaborators with low latency.
Each browser tab that has opened the shared file is a REPLICA. And note: even if the app does NOT allow offline editing, the fact that multiple users can make edits without waiting for a server response ALREADY MAKES IT MULTI-LEADER.
Both offline editing and real-time collaboration need the same infrastructure: capture the user's changes and either send them immediately (online) or store locally for later (offline); receive changes from collaborators, merge them into the local copy, and update the UI. If multiple users changed the file concurrently, conflict resolution logic is needed.
The vocabulary:
| Term | Meaning |
|---|---|
| Sync engine | A software library supporting this process. The idea is old; the term has recently gained attention |
| Offline-first | An app that allows the user to continue editing while offline |
| Local-first | Collaborative apps that are not only offline-first but designed to keep working even if the DEVELOPER SHUTS DOWN ALL THEIR ONLINE SERVICES — achieved with a sync engine using an open standard sync protocol with multiple available providers. Git is a local-first collaboration system (albeit without real-time collaboration) — you can sync via GitHub, GitLab, or any other host |
| Netcode | The multiplayer-video-game equivalent. Techniques are quite specific to games and don't directly carry over |
Four advantages of sync engines (vs the dominant model of keeping little state on the client and hitting a server for everything):
- Speed. Local data means the UI responds much faster. Some apps aim to respond to input in the NEXT FRAME — rendering within 16 ms on a 60 Hz display.
- Offline. An app doesn't need a separate offline mode: being offline is the same as having a very large network delay.
- Simpler programming model. Every service call requires error handling (Ch 5 §7.3) — if an update request fails, the UI must reflect that error. A sync engine lets the app read and write LOCAL data; these operations almost never fail, leading to a more DECLARATIVE programming style.
- Real-time updates. A sync engine combined with a reactive programming model is a good way to receive edit notifications and update the UI.
The limitation:
Sync engines work best when ALL the data the user may need is downloaded in advance and stored persistently on the client. Fine for all the files a user created (one user doesn't generate that much data); not suitable for the entire catalog of an ecommerce website.
History and implementations: pioneered by Lotus Notes in the 1980s (without the term). Today: proprietary backends (Google Firestore, Realm, Ditto) and open source backends suitable for local-first software (PouchDB/CouchDB, Automerge, Yjs).
3.4 Dealing with conflicting writes
The biggest problem with multi-leader replication — both in a geo-distributed server-side database and a local-first sync engine — is that concurrent writes on different leaders lead to conflicts that need to be resolved.
USER 1 USER 2
sets title A → B sets title A → C
│ │
▼ ▼
[LEADER 1] ═══════ async replication ═══════ [LEADER 2]
│ │
└──────────────▶ CONFLICT DETECTED ◀─────────┘
(This problem does not occur in a single-leader database.)Definition of concurrent: the two writes are concurrent because NEITHER WAS "AWARE" OF THE OTHER at the time it was originally made. It doesn't matter whether they literally happened at the same time — if made while offline, they might have happened some time apart. What matters is whether one write occurred in a state where the other had already taken effect.
Strategy 1 — Conflict avoidance
Ensure all writes for a particular record go through the same leader. Then conflicts cannot occur even if the database as a whole is multi-leader.
- Not possible for a sync engine client being updated offline, but sometimes possible in geo-replicated server systems
- Example: in an app where a user can edit only their own data, route requests from a particular user always to the same region. Different users have different "home" regions (picked by geographic proximity); from any one user's point of view, the configuration is essentially single-leader
- Another example: odd/even autoincrement — with two leaders, one generates only odd numbers, the other only even, so they can't concurrently assign the same ID
⚠️ Conflict avoidance BREAKS DOWN if you allow the leader to be changed. You might want to change the designated leader — a region is unavailable, or a user moved closer to a different region — and there is now a risk the user performs a write WHILE THE CHANGE IS IN PROGRESS, leading to a conflict.
Strategy 2 — Last write wins (LWW)
Attach a timestamp to each write; always use the value with the greatest timestamp. Ties broken by comparing the values (e.g. alphabetically earliest string).
The term is misleading: when two writes are CONCURRENT, which one is most recent is UNDEFINED, so the timestamp order of concurrent writes is essentially RANDOM.
The real meaning of LWW: when the same record is concurrently written on different leaders, ONE OF THOSE WRITES IS RANDOMLY CHOSEN AS THE WINNER AND THE OTHERS ARE SILENTLY DISCARDED, even though they were successfully processed by their respective leaders. This achieves eventual consistency AT THE COST OF DATA LOSS.
When LWW is fine: if you can avoid conflicts — e.g. only inserting records with a unique key and never updating them. If you update existing records, or different leaders may insert records with the same key, you have to decide whether lost updates are acceptable.
The clock hazard: with a real-time clock (Unix timestamp), the system becomes very sensitive to clock synchronization. If one node's clock is AHEAD of the others, your attempt to overwrite a value written by that node MAY BE IGNORED because it has a lower timestamp — even though it clearly occurred later. Solvable with a logical clock (Ch 10).
Strategy 3 — Manual conflict resolution
Like a Git merge conflict — but it would be impractical for a conflict to stop the entire replication process until a human resolves it. Instead databases store all concurrently written values (siblings) and return all of them on the next read; you resolve them automatically in application code (e.g. concatenate B and C into B/C) or by asking the user, then write back a new value. Used by CouchDB.
Four problems:
- The API of the database changes — a title that was a string becomes a set of strings usually containing one element. Awkward in application code.
- Asking the user to merge is a lot of work — for the developer (building conflict-resolution UI) and the user (who may be confused about what they're being asked and why). In many cases it's better to merge automatically than to bother the user.
- Naive automatic merging is surprising. The Amazon shopping-cart anomaly:
START: cart = {Book, DVD, Soap}
DEVICE 1 removes Book → {DVD, Soap} ┐
├── MERGE BY SET UNION
DEVICE 2 removes DVD → {Book, Soap} ┘ ▼
{Book, DVD, Soap}
▲ ▲
BOTH REMOVED ITEMS REAPPEAR IN THE CART- Resolution can itself introduce a new conflict. If multiple nodes observe and concurrently resolve the same conflict, one node may merge B and C into
B/Cand another intoC/B. When that conflict is merged, you may getB/C/C/Bor something similarly surprising.
Strategy 4 — Automatic conflict resolution
Automatic conflict resolution ensures all replicas CONVERGE to the same state — all replicas that have processed the same set of writes have the same state, REGARDLESS OF THE ORDER in which the writes arrived. Combining eventual consistency with a convergence guarantee is STRONG EVENTUAL CONSISTENCY.
LWW is the simplest such algorithm. More sophisticated merge algorithms exist per data type, with the goal of preserving the intended effect of all updates as much as possible, and hence avoiding data loss:
| Data type | Merge approach |
|---|---|
| Text (wiki title/body) | Detect which characters were inserted or deleted between versions; preserve all insertions and deletions made in any sibling. Concurrent insertions at the same position are ordered deterministically so all nodes get the same outcome |
| Collection (ordered to-do list, unordered cart) | Merge like text, tracking insertions and deletions. This is what fixes the Amazon anomaly: the algorithms track that Book and DVD were DELETED, so the merged result is {Soap} |
| Counter (likes on a post) | Tell how many increments and decrements happened on each sibling and add them together, so the result does not double-count and does not drop updates |
| Key-value map | Merge updates to the same key by applying one of the other algorithms to the values; updates to different keys are handled independently |
The limit: if you want to enforce that a list contains no more than five items, and multiple users concurrently add so there are more than five, your only option is to DROP SOME OF THE ITEMS.
Nevertheless, automatic conflict resolution is sufficient to build many useful apps. And if you start from the requirement of building a collaborative offline-first or local-first app, CONFLICT RESOLUTION IS INEVITABLE, and automating it is often the best approach.
CRDTs vs OT
Both perform automatic merges for all the types above; different design philosophies and performance characteristics.
The worked example: two replicas start with ice. One prepends n → nice; the other concurrently appends ! → ice!. Both must converge to nice!.
| Operational transformation (OT) | CRDT |
|---|---|
Record the index of each operation: insert 'n' at index 0, insert '!' at index 3. | Give each character a unique, immutable ID: i=1A, c=2A, e=3A. |
Exchange operations. 'n' at index 0 applies as is, but '!' at index 3 applied to “nice” would give “nic!e” — wrong. | An insert carries the ID of the new character (4B) and the ID of the existing character after which to insert (3A). Inserting at the beginning passes NIL. |
So the index must be transformed to account for concurrent operations already applied: '!' → index 4. | Concurrent insertions at the same position are ordered by character ID, so replicas converge without any transformation. |
Where each is used: OT is most often used for real-time collaborative text editing (Google Docs). CRDTs are found in distributed databases (Redis Enterprise, Riak, Azure Cosmos DB). Sync engines for JSON can be implemented with either — CRDTs (Automerge, Yjs) or OT (ShareDB). It's possible to combine the advantages of both in one algorithm.
Lists and arrays work the same way with list elements instead of characters; other datatypes like key-value maps can be added quite easily.
3.5 Types of conflict — the subtle kind
Obvious conflict: two writes concurrently modified the same field of the same record.
Subtle conflict — the meeting room booking system:
The system inserts a NEW RECORD for each booking rather than updating a field. The application must ensure each room is booked by only one group at any one time — no overlapping bookings.
A conflict arises if two bookings are created for the same room at the same time. EVEN IF THE APPLICATION CHECKS AVAILABILITY BEFORE ALLOWING A BOOKING, a conflict can arise if the two bookings are made close enough together that both see the room as unbooked prior to inserting their record.
There isn't a quick ready-made answer. More examples in Ch 8 (write skew / phantoms); scalable approaches to detecting and resolving in Ch 13.
6.2 Problems with Replication Lag
And on where to solve it: dealing with these issues in application code is complex and easy to get wrong.
6.4 Leaderless Replication
In some implementations the client sends writes to several replicas directly; in others a coordinator node does this on the client's behalf.