skip to content

What is the difference between primary-replica (leader-follower) replication and multi-primary (multi-leader) replication in a distributed database?

level: juniorimportance: must knowfreq 70%

answer

  1. one writer many readers
  2. WAL/binlog streaming
  3. async vs semi-sync
  4. conflict resolution needed for multi-leader
  5. failover promotion

basics

~20 s

In primary-replica setups, one server accepts writes and copies them to other servers that only serve reads. In multi-primary, several servers can all accept writes and have to sync changes with each other, which can cause conflicts.

solid answer

~40 s

Primary-replica replication designates a single node as the primary that accepts all writes; it streams a change log (e.g. a write-ahead log) to one or more replicas which apply changes and serve read traffic, giving a simple, conflict-free consistency model but a single write bottleneck and a failover gap. Multi-primary replication allows two or more nodes to accept writes independently, which improves write availability and lets writes land closer to users in multi-region setups, but requires conflict detection/resolution (last-write-wins, vector clocks, CRDTs, or application-level merge) because concurrent writes to the same record on different primaries can diverge. Most systems default to primary-replica for its operational simplicity and only adopt multi-primary when write locality or write availability during a partition is a hard requirement.

go deeper

for a junior

Can state the basic distinction and that replicas typically serve reads only.

for a middle

Understands async vs sync trade-offs and can describe the mechanics of failover.

for a senior

Can pick an appropriate replication topology for given consistency/availability requirements and knows concrete conflict-resolution strategies.

for a principal

Can architect a multi-region replication topology weighing write locality, RPO/RTO, and select real systems/products accordingly.

## What replication buys you Replication means keeping copies of the same data on multiple machines so a system can survive node failure and spread out read traffic. The two dominant topologies for organizing who is allowed to write are **primary-replica** (leader-follower, master-slave) and **multi-primary** (multi-leader, or leaderless in its most decentralized form). ## Primary-replica: one node owns the writes In primary-replica replication, exactly one node - the **primary** - is the source of truth for writes. Every write goes through the primary, which appends the change to a durable log (a write-ahead log in Postgres, a binlog in MySQL, a commit log in many NoSQL stores). That log is then streamed to one or more replicas, which apply the entries in order and expose a read-only copy of the data. Reads can be routed to the primary for freshness or to replicas to spread load, at the cost of potentially reading stale data. Replication can be: | Mode | What the primary does before it confirms | |---|---| | **Synchronous** | waits for at least one replica to acknowledge before confirming the write to the client | | **Semi-synchronous** | waits for acknowledgment but not application | | **Fully asynchronous** | confirms immediately and ships the log in the background | ## Multi-primary: several nodes accept writes Multi-primary replication removes the single-writer constraint: two or more nodes each accept writes directly and then propagate those writes to their peers. This is attractive when you need writes to succeed even during a network partition that isolates one region, or when you want writes to land in the data center closest to the user to cut latency - e.g. a multinational app that lets users in Tokyo and Frankfurt both write to a 'local' primary. The catch is that two primaries can accept conflicting writes to the same record before either has heard from the other. The system must then detect and resolve that conflict: - **last-write-wins** - compare timestamps and discard the loser - simple but silently loses data; - **vector clocks or version vectors** - detect that two writes are concurrent and surface both to the application or a merge function; - **CRDTs** - data types engineered so concurrent updates merge deterministically without loss, e.g. a grow-only counter. ## The core trade-off The core trade-off is operational simplicity and consistency versus write availability and locality. - **Primary-replica** is easier to reason about - there is one linear history of writes - and is the default choice for the large majority of relational deployments (Postgres streaming replication, MySQL binlog replication, Amazon Aurora's single-writer model). Its weakness is that the primary is both a write bottleneck and a single point of failure: if it goes down, writes stop until a replica is promoted, and that promotion is not instantaneous. - **Multi-primary** buys continuous write availability and can reduce write latency for geographically distributed users, but every additional primary is another source of conflicting writes, and conflict-resolution logic is genuinely hard to get right and to test. ## Failure modes Failure modes differ accordingly. - In primary-replica systems, the classic failure is **split-brain**: a network partition makes the cluster believe the old primary is dead, promotes a replica, and then the old primary - still alive but isolated - keeps accepting writes it thinks are valid. Two 'primaries' now exist simultaneously, and reconciling their histories afterward can mean silent data loss. Modern systems avoid this with consensus-based leader election (Raft or Paxos, as used by etcd/Patroni-managed Postgres, or CockroachDB's per-range Raft groups) and fencing tokens that let a new primary invalidate the old one's ability to write. - Asynchronous replication also has a quieter failure mode: **replication lag**. If the primary fails before its latest writes reach any replica, those writes are lost on failover - this is the classic availability/durability trade governed by RPO (recovery point objective). - In multi-primary systems, the equivalent failure is a conflict that the resolution policy handles badly: last-write-wins can silently drop a legitimate update if clocks are skewed, and merge functions that are not truly commutative/associative can produce different results on different nodes, breaking convergence. ## Systems you can name A concrete example: - **MySQL Group Replication and Galera Cluster** support genuine multi-primary writes with certification-based conflict detection (a write is 'certified' cluster-wide before committing, and conflicting transactions are aborted rather than silently merged). - **Amazon DynamoDB and Cassandra** are leaderless - any replica can accept a write - and rely on quorum reads/writes plus last-write-wins or client-supplied conflict resolution. - **Google Spanner and CockroachDB**, by contrast, give the appearance of a single global database while internally running per-shard Raft groups, each with its own leader - so at the shard level it is still primary-replica, just partitioned finely enough that no single primary is a global bottleneck.

  • In a primary-replica setup, why can a client that just wrote data sometimes not see that data on its next read?
    Because replication is usually asynchronous, so there's a short window (replication lag) where a replica hasn't yet applied the write. The fix is to route that user's next read to the primary, or to a replica confirmed to have caught up, rather than a random replica.
  • What is split-brain in a primary-replica cluster, and how do systems prevent it?
    Split-brain is when a network partition causes two nodes to both believe they are the primary and both accept writes, producing diverging histories. Systems prevent it with consensus-based leader election (Raft/Paxos) and fencing tokens that revoke the old primary's ability to write once a new one is elected.
  • Name a real system that uses multi-primary replication and how it resolves write conflicts.
    Galera Cluster (MySQL) uses certification-based replication: a transaction is certified across the cluster before commit, and conflicting concurrent transactions are aborted rather than silently merged. Cassandra and DynamoDB, which are leaderless, typically use timestamp-based last-write-wins or expose version conflicts to the client.

Primary-replica is like a single teacher writing on the whiteboard while students copy it into their notebooks - fast and consistent, but if the teacher is out sick, nothing new gets written until someone takes over. Multi-primary is like several teachers writing on shared whiteboards in different classrooms and phoning each other to reconcile any contradictions.

saying these in an interview costs you the question

  • says replicas can always be written to in a primary-replica setup
  • doesn't mention conflict resolution when describing multi-primary
  • thinks failover is instantaneous with zero data loss by default
  • conflates replication with sharding
  • assumes synchronous replication has no latency cost

context