Replication and Consistency in Distributed Systems: Single-Leader, Multi-Leader, and Leaderless Topologies
Introduction
Every system design interview eventually reaches the same moment. You have drawn a database box, the interviewer points at it, and asks: "What happens when this node dies?" The honest answer is "we keep copies," and the real conversation starts right after: how many copies, who is allowed to write to them, how quickly they agree, and what a reader sees while they disagree.
That is replication, and its twin is consistency: the contract a replicated system makes about which values a read can return. Candidates often know the vocabulary (leader, quorum, eventual consistency, CAP) but blur the definitions, and interviewers notice. "Cassandra is AP" and "quorums are strongly consistent" both sound right and are both imprecise in ways that matter.
This post builds a precise mental model of the three replication topologies (single-leader, multi-leader, and leaderless) using one running example: a global e-commerce store with users in Mumbai and Northern Virginia, where user 42 keeps a cart, updates a shipping address, and places orders. For each topology we cover when to use it, how it works, what user 42 experiences, how it scales, and how it fails. Then we map the consistency spectrum, pin down CAP and PACELC, and look at how PostgreSQL, DynamoDB, Cassandra, and Spanner choose.

Why We Replicate
Replication means keeping a copy of the same data on several machines connected by a network. Teams do it for four distinct reasons, and naming the reason up front shapes every later decision.
- Durability. Disks and machines get lost. Copies in other failure domains (racks, availability zones, regions) keep user 42's order history alive.
- Availability. If one node is down for a crash, a patch, or a deploy, another keeps serving.
- Read scaling. Replicas spread read traffic, which is why "add read replicas" is the first scaling move for a read-heavy relational workload.
- Latency. A user in Mumbai reading from a Mumbai replica avoids the round trip to Virginia.
Replication is not partitioning (sharding). Partitioning splits data so each node holds a subset; replication copies each subset to several nodes. Real systems do both. In an interview, "we shard by user ID" answers a capacity question, while "each shard has three replicas across zones" answers a failure question.
The hard part is handling change: a write lands on one machine first, and until every copy has it, the copies disagree. Every topology below is a different answer to two questions: which nodes may accept a write, and when does the system tell the client the write succeeded?
Single-Leader Replication: Sync vs Async and Failover
When to Use It
Single-leader (also called leader-follower or primary-replica) replication is the default for PostgreSQL, MySQL, MongoDB replica sets, and the per-partition replication inside many distributed stores. Use it when you need one unambiguous order of writes: checkout, payments, inventory decrements, uniqueness constraints. One node decides the order, so there are no write conflicts.
How It Works
All writes go to the leader, which applies them and appends them to a replication log (the WAL in PostgreSQL, the binary log in MySQL, the oplog in MongoDB). Followers stream the log and apply changes in the same order. Reads can go to the leader (current) or to followers (possibly behind).
The critical knob is when the leader acknowledges the write to the client:
- Asynchronous replication. The leader commits locally and replies immediately; followers catch up later. Write latency is lowest, and a slow follower never blocks writes. The cost: if the leader dies before a follower receives the latest changes, those acknowledged writes are lost on failover. Asynchronous is the default for PostgreSQL streaming replication and for standard MySQL replication.
- Synchronous replication. The leader waits until at least one follower confirms before replying. An acknowledged write now exists on two machines. The cost is latency (every commit includes a network round trip) and availability: if the synchronous follower is down and no other qualifies, writes stall.
- Semi-synchronous replication. The usual compromise: wait for one follower out of several. In MySQL, the leader waits until at least one replica has received the transaction and written it to its relay log. If none acknowledges within a configurable timeout, MySQL falls back to asynchronous mode, silently dropping the durability guarantee until a replica catches up. Semi-sync protects you only while it is actually in effect.
PostgreSQL expresses the same idea through synchronous_standby_names and synchronous_commit, applied to user 42's checkout:
# postgresql.conf on the leader
# Wait for any one of two named standbys (quorum syntax, PostgreSQL 10+)
synchronous_standby_names = 'ANY 1 (replica_az_b, replica_az_c)'
# on: wait until the standby has flushed the WAL to disk
# remote_write: wait until the standby has written it (not yet fsynced)
# remote_apply: wait until the standby has applied it, so reads there see it
synchronous_commit = on
synchronous_commit can also be set per transaction, so a team can keep on for orders and use SET LOCAL synchronous_commit = off for low-value writes such as page views.
Replication Lag and Its Anomalies
With asynchronous followers, there is a window (usually milliseconds, sometimes seconds or minutes under load or during a long-running transaction) where followers lag behind the leader. If the application reads from followers, users can observe three well-known anomalies. Here is what user 42 might experience:
Read-your-writes violation t0 user 42 changes shipping address to "Bandra" -> leader t1 page reloads, reads profile -> follower F1 (lagging) t1 sees old address "Andheri"; thinks the save failed and retries Monotonic reads violation t0 order #881 status becomes "shipped" -> leader t1 user 42 views order list -> follower F1 (caught up): "shipped" t2 refreshes -> follower F2 (lagging): "processing" status appears to move backwards in time Consistent prefix violation (most visible across partitions) t0 support agent writes: "Refund issued" -> partition P1 t1 user 42 replies: "Thanks, I see it" -> partition P2 a reader whose P2 replica is fresher than its P1 replica sees the thank-you before the refund message it is replying to
The fixes are targeted. For read-your-writes, read from the leader for a short window after a user writes, or have the client remember its last write's log position (a PostgreSQL LSN, for example) and read only from a replica that has replayed past it. For monotonic reads, pin each user to one replica by hashing the user ID. For consistent prefix, keep causally related writes in the same partition. Read-your-writes and monotonic reads, plus monotonic writes and writes-follow-reads, are the four session guarantees described by Terry and colleagues in 1994; consistent prefix comes from later work.
Failover and Split Brain
When the leader fails, a follower must be promoted. Failover has three steps: detect the failure (usually a timeout on heartbeats), choose a new leader (ideally the most up-to-date follower), and reconfigure clients and other followers to use it. Each step can go wrong.
Detection is a guess: a leader stalled by a garbage collection pause or a flaky link looks exactly like a dead one. Too low a timeout causes needless failovers during load spikes; too high a timeout means longer outages.
The dangerous case is split brain: the old leader is only unreachable, not dead, and keeps accepting writes while a new leader does too. Two leaders means two diverging histories, and a single-leader design has no way to merge them.
The robust way to choose a leader is to delegate the decision to consensus. Raft and Paxos solve the same core problem: getting nodes to agree on a value (here, "who leads term N") despite failures and delayed messages. Conceptually, Raft works like this:
- Time is divided into numbered terms, each with at most one leader.
- A follower that stops hearing heartbeats waits a randomized timeout, increments the term, and requests votes.
- A node votes at most once per term, and only for a candidate whose log is at least as up to date as its own. Winning needs a majority.
- Any two majorities overlap, so two leaders cannot win the same term, and a new leader always holds every committed entry.
With five nodes, a Raft cluster tolerates two failures; with three, it tolerates one. Production PostgreSQL clusters commonly get this property indirectly: PostgreSQL has no built-in automatic failover, so tools such as Patroni hold a leader lease in a consensus-backed store such as etcd, ZooKeeper, or Consul, and only the lease holder runs as leader.
Consensus prevents two leaders from being elected, but not an old leader that wakes from a pause unaware it was replaced. The defence is a fencing token: each lease grant comes with a monotonically increasing number (an etcd revision or a ZooKeeper transaction ID works). The leader attaches it to every write, and storage rejects any token lower than the highest it has seen.

Where the storage cannot check tokens, operators fall back to cruder fencing: cut the old leader's network access, revoke its virtual IP, or power it off (often called STONITH, "shoot the other node in the head").
Scaling
Single-leader scales reads by adding followers, including remote ones for local reads. It does not scale writes beyond one node per partition. The escape is to partition and run single-leader replication per partition, as MongoDB sharded clusters, Vitess on MySQL, and CockroachDB's per-range Raft groups do.
Failure Modes and Anti-Patterns
- Asynchronous failover loses acknowledged writes. Promote a follower that was 200 ms behind and the last 200 ms of commits vanish. If those order IDs were already cached or sent to a payment provider, the new leader can reissue them to different orders, a well-known way for data to leak between users.
- Reading from followers by default. Fine for product listings, wrong for "show me the order I just placed." Decide per query.
- Treating semi-sync as always-on. Monitor for the fallback to asynchronous mode, or you will discover it during the incident where it mattered.
Multi-Leader Replication: Conflicts, Last-Writer-Wins, and CRDTs
When to Use It
Multi-leader (also called active-active or multi-primary) replication lets more than one node accept writes, each replicating asynchronously to the others. Three situations justify the complexity:
- Multi-region writes with local latency. User 42 in Mumbai writes to a Mumbai leader; a user in Virginia writes to a Virginia leader. Neither waits for a cross-ocean round trip, and each region keeps accepting writes if the link between them fails.
- Offline-capable clients. An app that works offline is a leader with a local database that syncs later, the model CouchDB is built around.
- Collaborative editing. Each participant edits a local replica and merges others' changes.
If none of these apply, prefer single-leader. Multi-leader trades away the single write order, and with it you inherit conflicts.
How It Works and Why Conflicts Happen
Each leader applies writes locally, acknowledges immediately, and ships its changes to the other leaders in the background. A conflict occurs when two leaders accept concurrent writes to the same data before seeing each other's change.
"Concurrent" means neither write knew about the other; it is not about wall-clock time. If user 42 adds a book on a phone connected to Mumbai, then a second later removes a mug on a laptop routed to Virginia, and the regions had not exchanged changes in between, the writes are concurrent.
To detect concurrency, systems track causality with version vectors (a close relative of vector clocks, which track per-process event counts). Each replica keeps a counter per leader, and every version of a value carries the vector of counters it was derived from:
cart:42 starts as {mug} version vector {mumbai: 3, virginia: 5} Mumbai: add "book" -> {mug, book} {mumbai: 4, virginia: 5} Virginia: remove "mug" -> {} {mumbai: 3, virginia: 6} Compare {4,5} with {3,6}: mumbai counter: 4 > 3 virginia counter: 5 < 6 neither vector dominates the other, so the writes are concurrent: CONFLICT If Virginia had first received Mumbai's change, its version would be {mumbai: 4, virginia: 6}, which dominates {4,5}: a normal overwrite, no conflict.
Once a conflict is detected, something must resolve it.

Resolution Strategy 1: Last-Writer-Wins
Last-writer-wins (LWW) keeps the write with the highest timestamp and discards the rest. Every replica picks the same winner, so it converges. DynamoDB global tables use last-writer-wins to reconcile concurrent updates in their multi-region eventually consistent mode, and Cassandra applies LWW per column using write timestamps.
The cost is data loss by design. In the cart example, LWW keeps either {mug, book} or {}, so one change silently disappears. Clocks drift, too, so a node whose clock runs ahead wins conflicts it should lose. LWW suits values where any one recent write is acceptable (a last-seen timestamp, a profile photo URL), not anything that accumulates.
Resolution Strategy 2: Application Merge
The system keeps all concurrent versions (siblings) and hands them to the application on the next read. The Amazon Dynamo paper (2007) did this for the shopping cart, merging by union. Its known flaw: a removed item can reappear, because a union cannot tell "never added" from "added then removed."
Resolution Strategy 3: CRDTs
Conflict-free replicated data types (CRDTs) are data structures designed so that concurrent updates merge automatically and deterministically, with all replicas converging to the same value regardless of the order in which they receive updates. Common building blocks:
- G-Counter and PN-Counter: each replica increments its own slot; the value is the sum. Good for like counts and view tallies (but not for enforcing "never below zero").
- OR-Set (observed-remove set): each add is tagged with a unique ID; a remove deletes only the tags it has observed. A concurrent add carries a tag the remove never saw, so it survives. For the cart, the merge of "add book" and "remove mug" correctly yields
{book}. - Sequence CRDTs (RGA and its successors): ordered lists of characters for collaborative text, used by libraries such as Automerge and Yjs.
CRDTs guarantee convergence, not business correctness. Two regions can each sell the last unit of a product, and a PN-Counter will happily converge to minus one. Invariants that span replicas ("stock never below zero," "username is unique") need coordination, which means routing those writes to a single leader or using consensus.
Scaling
Multi-leader scales writes geographically: each region absorbs its own writes. With many leaders, all-to-all replication is the most robust; ring and star topologies need fewer links but create single points of failure.
Failure Modes and Anti-Patterns
- Hidden conflicts in "low-conflict" data. Background jobs, retries, and admin tools edit the same rows users do.
- Auto-increment keys. Two leaders generating IDs from their own sequences will collide. Use UUIDs, or interleaved ranges per leader.
- Uniqueness across leaders. Two users can claim the same username in two regions.
- Causality violations. An update can arrive at a region before the insert it modifies. Version vectors detect this; naive timestamp ordering does not.
Leaderless Replication: Dynamo-Style Quorums and Read Repair
When to Use It
Leaderless replication, popularized by the Dynamo paper and implemented by Cassandra, Riak, and ScyllaDB, has the client (or a coordinator node) send each write to several replicas directly. With no leader, there is no failover pause. Use it for high write volumes that must stay writable through node failures and tolerate eventual consistency: carts, activity feeds, sensor readings, sessions.
How Quorums Work
Each key is stored on N replicas, usually chosen by consistent hashing. A write goes to all N and succeeds once W acknowledge. A read succeeds once R replicas respond, and the newest version wins. If W + R > N, the write set and the read set must overlap in at least one node, so the read sees the latest successful write.
For user 42's cart with N = 3, W = 2, R = 2:
PUT cart:42 = v2 -> sent to A, B, C A ack, B ack -> W=2 reached: success returned to client C is slow; it still holds v1 GET cart:42 -> sent to A, B, C; first two responses used responses: C = v1, A = v2 overlap guaranteed by W + R = 4 > N = 3, so v2 is among them client returns v2 (newest version) coordinator writes v2 back to C <- read repair
Common configurations in Cassandra terms: QUORUM for both reads and writes gives the overlap; ONE for both is fastest and weakest; LOCAL_QUORUM computes the quorum within the local datacenter only, which avoids cross-region latency but gives up overlap with writes made in the other region.

Keeping Replicas in Sync
With no ordered log to replay, three mechanisms bring lagging replicas back in line:
- Read repair. When a read sees a stale replica, the coordinator writes the newer value back to it. This repairs frequently read keys quickly and rarely read keys never. (Cassandra 4.0 removed the older probabilistic background read repair options; blocking read repair at quorum levels remains.)
- Hinted handoff. If a target replica is down during a write, another node stores a "hint" and replays it when the replica returns. In Cassandra, hints are kept for a bounded window (3 hours by default, set by
max_hint_window); a node down longer than that needs a full repair. - Anti-entropy with Merkle trees. Each replica builds a hash tree over its key ranges: leaves hash small ranges, parents hash their children. Two replicas compare roots and descend only into subtrees whose hashes differ, isolating the few ranges to stream. In Cassandra this is
nodetool repair, which must run at least once withingc_grace_seconds(10 days by default) so deletions propagate before tombstones are purged. Skip it, and deleted data can come back.
Sloppy Quorums
What if the client can reach only one of the three designated replicas for cart:42? A strict quorum fails the write. A sloppy quorum, described in the Dynamo paper, accepts it on any W reachable nodes, even outside the key's designated N, and uses hinted handoff to move the data home later. Write availability goes up, but the overlap guarantee breaks: a read quorum may miss every stand-in holding the latest write.
Implementations differ here. Riak exposes sloppy quorums. Cassandra stores hints but does not count them toward the consistency level (except for the special level ANY), so its quorums remain strict and a write simply fails if not enough designated replicas respond.
Quorums Are Not Linearizability
W + R > N is often described as "strong consistency." It is not, by itself, linearizability. Edge cases that break it:
- Two concurrent writes with LWW conflict resolution: one is silently dropped, and which one depends on clocks.
- A write that succeeds on fewer than W replicas is reported as failed but is not rolled back; later reads may or may not see it.
- A write in progress: one reader sees the new value from the replica that already has it, and a later reader hits two replicas that do not yet have it, observing time going backwards.
- Sloppy quorums, as above.
Readers that synchronously repair before returning close some of these gaps, but still give no atomic compare-and-set. For that, Cassandra offers lightweight transactions (IF NOT EXISTS, IF balance = ...) built on a Paxos round, at notably higher latency than a normal write.
Scaling
Adding nodes to the ring scales reads and writes; consistent hashing moves only a fraction of the data. A single hot key still hammers its N replicas.
Failure Modes and Anti-Patterns
- Assuming QUORUM means correct. It narrows windows; it does not give compare-and-set or uniqueness.
- Never running repair. Leads to zombie data and replicas that drift apart.
- Read-modify-write without versions. Two clients read a counter, increment, and write back; LWW keeps one increment. Use counters or CRDTs, or lightweight transactions.
The Consistency Spectrum: Linearizable to Eventual
"Consistency" is overloaded. The C in ACID means application invariants hold. The C in CAP means linearizability. Here we mean the replication sense: which values may a read return, given the writes that have happened? Models form a spectrum, where stronger models need more coordination and therefore more latency and less availability under failures.

Linearizability makes the system behave as if there were a single copy of the data, with each operation taking effect atomically at some instant between its start and end. Once any read returns a new value, every later read must return it or something newer. You need it for leader election, locks, uniqueness checks, and selling the last unit exactly once. etcd reads are linearizable by default; ZooKeeper writes are linearizable, but a read can hit a lagging server unless the client calls sync first.
Sequential consistency keeps a single total order of operations that all clients agree on and that respects each client's own program order, but that order need not match real time. A read can be stale as long as everyone sees the same history.
Causal consistency orders only causally related operations: if write B was made by someone who had seen write A, every node shows A before B. Concurrent writes may appear in different orders on different nodes. The refund message always precedes the reply to it. Causal consistency is among the strongest models that can stay available during a partition. MongoDB offers causally consistent sessions (since version 3.6) with majority read and write concerns.
Session guarantees are per-client promises: read-your-writes, monotonic reads, monotonic writes, and writes-follow-reads. They are weaker than causal consistency across clients but fix most user-visible anomalies, and they are cheap to implement with sticky routing or log-position tokens.
Eventual consistency promises only that if writes stop, replicas converge. It says nothing about what you read in the meantime. DynamoDB's default reads and Cassandra at ONE sit here.
Two clarifications earn points. First, serializability is a different axis: it is an isolation guarantee about multi-object transactions (the result equals some serial order), not about recency. A serializable system can serve stale reads. Strict serializability is serializability plus linearizability, and it is what Spanner calls external consistency. Second, the spectrum is not just a property of a database, but of how you use it: the same Cassandra cluster is eventually consistent at ONE and much closer to strong at QUORUM.
CAP and PACELC in Practice
What CAP Actually Says
The CAP theorem (conjectured by Eric Brewer in 2000, proved by Gilbert and Lynch in 2002) is narrower than its reputation:
- C is linearizability, not "consistency" in a loose sense, and not the C in ACID.
- A means every request received by a non-failing node gets a non-error response, with no time bound beyond "eventually."
- P is a network partition: messages between groups of nodes are lost.
The theorem says: while a network partition is happening, a system must choose between linearizability and availability. Nodes on one side cannot learn about writes on the other side, so either they refuse to answer (losing A) or they answer possibly stale (losing C).
- "Pick two of three" is misleading. You cannot opt out of partitions in a distributed system, so "CA" describes only a system that is not really distributed, or one that gives up availability as soon as a partition occurs. The real choice is C or A during a partition.
- CAP says nothing about normal operation. When the network is healthy, a system can be both linearizable and available.
- Most systems are neither strictly CP nor strictly AP. A single-leader database with asynchronous replicas is not linearizable even without a partition (replica reads are stale), and it is not available in the CAP sense either (nodes cut off from the leader cannot accept writes). Cassandra's position depends on the consistency level per query.
- CAP availability is not your SLA. A CP system with fast failover can still have excellent practical uptime.
PACELC: The Trade-Off You Face Every Day
Daniel Abadi's PACELC formulation (published in 2012) adds the missing half: if there is a Partition, choose Availability or Consistency; Else, in normal operation, choose Latency or Consistency. The everyday cost of strong consistency is not partition behavior, which is rare. It is the extra round trips on every operation to coordinate replicas.
Abadi classified Dynamo-style systems such as Cassandra and Riak (in their default configurations) as PA/EL: they stay available in partitions and favor low latency normally. Fully consistent systems that coordinate on every write are PC/EC. For user 42, this is concrete: a synchronous cross-region write from Mumbai to Virginia adds a cross-ocean round trip to every checkout, all the time, partition or not. That latency, not CAP, is usually what drives the design.
How Real Systems Choose: PostgreSQL, DynamoDB, Cassandra, Spanner
PostgreSQL
Single-leader. Streaming replication ships WAL to standbys, asynchronous by default, with synchronous and quorum modes via synchronous_standby_names (multiple synchronous standbys arrived in version 9.6, and the FIRST and ANY keywords in version 10). Logical replication (built in since version 10) replicates selected tables. There is no built-in automatic failover; teams use tools like Patroni. Standby reads are stale by the replication lag; synchronous_commit = remote_apply gives read-your-writes only on the standbys acting as synchronous at commit time, while asynchronous standbys still lag. Amazon Aurora keeps a single writer but, per its 2017 SIGMOD paper, stores six copies across three availability zones with a write quorum of four and a read quorum of three: a quorum system under a single-leader database.
DynamoDB
DynamoDB shares a name with the Dynamo paper but not its design. Per its 2022 USENIX ATC paper, each partition is replicated across three availability zones and uses Multi-Paxos to elect a leader; writes go through the leader and are acknowledged once a quorum of replicas persists them. Reads are eventually consistent by default; with ConsistentRead=true the leader serves the latest acknowledged write at twice the read capacity cost. Across regions, global tables have traditionally been multi-leader with last-writer-wins. AWS has since added a multi-region strong consistency option (announced in late 2024 and generally available in 2025); check current documentation for its constraints before relying on it.
Cassandra
Leaderless, with tunable consistency per query (ONE, QUORUM, LOCAL_QUORUM, EACH_QUORUM, ALL, and SERIAL variants for lightweight transactions). Conflict resolution is last-writer-wins per cell using write timestamps; Cassandra does not use vector clocks. Replica convergence relies on read repair, hinted handoff, and Merkle-tree repair. A common multi-region setup uses LOCAL_QUORUM for both reads and writes, giving read-your-writes for clients that stay in one region and eventual consistency across regions.
Spanner
Spanner splits data into shards and replicates each with Paxos, one leader per group. Its distinctive piece is TrueTime: an API returning the current time as an interval [earliest, latest] with bounded uncertainty, backed by GPS receivers and atomic clocks. At commit, Spanner assigns a timestamp and performs a commit wait, delaying the acknowledgment until that timestamp is certainly in the past everywhere. So if T1 commits before T2 starts, T1 has the smaller timestamp: external consistency (strict serializability) globally, plus lock-free snapshot reads. The price is write latency. Brewer's 2017 paper on Spanner calls it technically CP, but experienced as highly available because partitions are rare on Google's network.
Head-to-Head Comparison
| Dimension | Single-leader (async) | Single-leader (sync / consensus) | Multi-leader | Leaderless (quorum) |
|---|---|---|---|---|
| Who accepts writes | One leader per partition | One leader per partition | One leader per region or device | Any replica (coordinator fans out) |
| Write latency | Lowest (local commit) | Local plus one replica round trip | Local | Wait for W replicas |
| Data loss on node failure | Possible (unreplicated tail) | None for acknowledged writes | Possible until replicated | None if W or more replicas persisted |
| Write conflicts | None | None | Yes: LWW, merge, or CRDTs | Yes: LWW or siblings |
| Read consistency | Leader: strong; replicas: stale | Leader: linearizable (with leases or read index) | Eventual across leaders | Tunable; not linearizable by default |
| Behavior during partition | Side without the leader cannot write | Minority side cannot write | All sides keep writing | Writes succeed where W reachable |
| Failover | Promotion, risk of lost writes | Consensus election, seconds | None needed per region | None needed |
| Cross-region writes | High latency to remote leader | High latency, every write | Local latency | Local with LOCAL_QUORUM |
| Global invariants (uniqueness, stock) | Easy | Easy | Hard without extra coordination | Needs Paxos-based LWT |
| Operational complexity | Low | Medium | High (conflict handling) | Medium to high (repair, tuning) |
| Typical systems | PostgreSQL, MySQL replicas | Patroni clusters, etcd, Spanner, CockroachDB | DynamoDB global tables, CouchDB, CRDT apps | Cassandra, Riak, ScyllaDB |
When to Pick Which
Scenario 1: Checkout and Payments Ledger
User 42 places order #881 and is charged; a lost or duplicated charge is unacceptable. Choose single-leader with synchronous or semi-synchronous replication to a standby in another availability zone, with consensus-based failover (Patroni on etcd, or a managed equivalent). Read the user's own orders from the leader, or from a replica past their last write LSN. Serve catalog browsing from asynchronous replicas. If writes outgrow one leader, shard by a key that keeps each checkout inside one shard.
Scenario 2: Shopping Cart and Activity Data at Global Scale
The cart is written constantly, must stay writable through node or zone failures, and tolerates a stale view. Choose leaderless replication (Cassandra with LOCAL_QUORUM per region), or a managed eventually consistent option such as DynamoDB global tables (leader-based within a region, multi-leader across regions). Store each cart line as its own row keyed by (user_id, item_id), so edits to different items never overwrite each other, instead of one cart blob under LWW. At checkout, copy the cart into the ledger from Scenario 1. Run scheduled repair.
Scenario 3: Global Inventory or Account Balances Across Regions
The store now sells limited-edition items in both regions, and overselling is unacceptable. CRDTs converge but cannot stop two regions from each selling the last unit. Choose consensus-replicated distributed SQL (Spanner, or a similar system such as CockroachDB), so each inventory row has one Paxos or Raft leader and every decrement is linearizable. Place each row's leader near the region that sells it most, and keep data without the invariant (reviews, carts) on cheaper stores. For collaborative editing of a shared wishlist, by contrast, multi-leader with sequence CRDTs fits: there is no hard invariant, and offline editing matters more than global order.
In the Interview
How to Justify a Choice
Use the same four steps every time:
- State the invariant or tolerance. "Two buyers must never get the last unit," or "a cart that is a few seconds stale is acceptable."
- Name the guarantee that protects it. "That is a linearizable compare-and-set on the stock row," or "read-your-writes for the user's own session is enough."
- Pick the topology that delivers it at the lowest cost. "Single leader per item with consensus failover," or "leaderless with
LOCAL_QUORUM." - Name what you are giving up. "Buyers in the other region pay a cross-region round trip on checkout," or "a region that is cut off keeps selling carts but cannot see the other region's edits until the link heals."
Interviewers listen hardest for step four.
Common Traps
- "CAP means pick two." Partitions are not optional. Say "during a partition, this system chooses consistency, so the minority side rejects writes."
- Using C loosely. CAP's C is linearizability. ACID's C is invariants. Replication consistency is a spectrum. Say which one you mean.
- "Quorums are strongly consistent." W + R > N gives overlap, not linearizability, and does nothing for compare-and-set.
- "Cassandra is AP, Spanner is CP, done." Cassandra's behavior is per query; Spanner's real cost shows up as latency in normal operation, which is the PACELC half.
- "DynamoDB is Dynamo." DynamoDB uses leader-based replication per partition with Multi-Paxos; the leaderless design is from the earlier Dynamo paper.
- "Synchronous replication means no data loss, so failover is safe." Only if you fail over to the synchronous replica, and only while semi-sync has not fallen back to async.
- Ignoring split brain. Saying "a follower takes over" without mentioning how the old leader is fenced is half an answer.
Likely Follow-Up Questions
- "A user updates their profile and immediately sees the old value. Why, and how do you fix it?" Replica lag violating read-your-writes. Read from the leader for a short window after a write, or track the write's log position and read from a replica that has passed it.
- "How do you prevent two leaders after a network blip?" Elect through a consensus store with leases, and attach fencing tokens to writes so storage rejects the stale leader.
- "A Cassandra QUORUM write with one of three replicas down?" It succeeds on two acknowledgments; a hint for the third is replayed if it returns within the hint window, otherwise repair fixes it.
- "Why not just use LWW everywhere in multi-region?" It silently discards concurrent writes and depends on clock accuracy. Use it only where one surviving write is acceptable; otherwise use CRDTs or route the write to one leader.
- "How does Spanner give global strong consistency?" Paxos per split for replication, plus TrueTime and commit wait so timestamp order matches real-time order. It pays with write latency.
- "Your async replica is the only survivor after the leader's disk dies. What do you tell the business?" That writes acknowledged after the replica's last received position are lost, how large that window was (from lag metrics), and that the fix is a synchronous standby for this data.
Related reading: Database Types Compared, Designing a Coordination Service like ZooKeeper or etcd, Globally Distributed SQL Databases like Spanner and CockroachDB, and Designing a Collaborative Editor like Google Docs.