Designing Data-Intensive Applications
Ch. 5

Multi-Leader Replication

Multiple leaders accept writes in different datacenters — but conflict resolution becomes necessary.

Single-leader replication means cross-datacenter writes traverse WAN latency. Multi-leader allows each site to write locally, then replicate. The cost is write conflicts when two leaders modify the same data concurrently.

In practice

CouchDB replicates bidirectionally between edge nodes and cloud — ideal for offline-first mobile apps. Cassandra's multi-datacenter replication lets each region write locally. Figma-style collaborative editors use CRDTs so concurrent edits merge without a single leader bottleneck.

Canva at scale

Collaborative design editing lets multiple users write concurrently from different regions. CRDTs and operational transforms merge conflicting layer edits without funneling every keystroke through one global leader.

typescript — Real-time collaboration over WebSocket
// Canva-style real-time collaboration over WebSocket
const wss = new WebSocketServer({ port: 8080 });

wss.on("connection", (ws, req) => {
  const designId = new URL(req.url!, "http://x").searchParams.get("designId");
  ws.on("message", (raw) => {
    const op = JSON.parse(raw.toString()) as DesignOperation;
    // Apply CRDT / OT merge, then broadcast to all peers on this design
    const merged = collabEngine.apply(designId!, op);
    broadcast(designId!, { type: "sync", state: merged });
  });
});
Diagram
Key Takeaways
  • Each datacenter has a local leader for low-latency writes.
  • Leaders replicate to each other asynchronously.
  • Concurrent writes to the same record on different leaders create conflicts.
  • Conflict resolution: last-write-wins, custom merge, or conflict-free replicated data types.
  • Multi-leader suits multi-datacenter deployments with offline tolerance.
  • CouchDB, Cassandra multi-DC, and CRDTs in collaborative apps use multi-leader patterns.
multi-leaderCassandraCouchDBCRDTwrite conflictlast-write-wins