In a single-leader (primary/replica) replication setup, how do writes and reads flow through the system, and why would a team adopt this design over just running one database server?
answer
- one writer, many readers
- WAL/binlog streamed to followers
- leader dies -> promote a follower
- lag causes stale reads
basics
~20 sOne server, the leader, is the only one allowed to accept writes. It records each change and sends a copy to the other servers (followers). Reads can be spread across the leader and followers. Teams do this to survive a server crash and to handle more read traffic than one machine could.
solid answer
~40 sSingle-leader replication elects one node as the leader; every write request is routed to it, appended to its write-ahead log, and streamed to one or more followers that apply the changes in the same order. Reads can be served by the leader for freshness or by followers for scale. Teams adopt it because it removes the ambiguity of concurrent writes to the same row from different nodes - there's one order of truth - while still getting horizontal read scaling, geographic read locality, and a standby that can be promoted if the leader fails. The cost is that all writes funnel through one node, and followers can lag behind.
go deeper
Should describe the basic shape: one leader takes writes, followers get copies, reads can go anywhere. Doesn't need failover mechanics.
Should articulate why reads can be stale (lag) and know that followers can be promoted on leader failure, even if they can't design the election protocol.
Should reason about the write-throughput ceiling this topology imposes, connect lag to specific anomalies (read-your-writes, monotonic reads), and know how to route reads to avoid them.
Should discuss this topology as a deliberate simplicity/throughput trade-off versus multi-leader and leaderless designs, and speak to failover risk (split-brain, data loss window) at an operational level across a fleet.
## What it is **Single-leader replication** (also called *primary-replica*, *master-slave*, or *active-passive* replication) is the most widely deployed replication topology because it sidesteps the hardest problem in replication: what happens when two nodes accept conflicting writes at the same time. ## The mechanism The mechanism is straightforward. 1. One node in the cluster is designated the **leader**. 2. Every write - insert, update, delete - is sent to the leader, which appends it to a durable, ordered log (often called a **write-ahead log** or **binlog**). 3. The leader then streams that log to every follower, and each follower replays the entries in the exact order they were written, so that, modulo lag, every follower's state is a delayed copy of the leader's state. 4. Reads can be routed to the leader (when you need the absolute latest data) or to any follower (when slightly stale data is acceptable, which is the common case for browsing, dashboards, or search). ## Why the design exists This design exists because coordinating writes across multiple independently-writable nodes is expensive and complex - it requires either locking across the network on every write or a mechanism to detect and merge conflicting writes after the fact. By funneling all writes through a single node, the leader naturally imposes a **total order** on every change, so followers never have to reconcile two different versions of the same row; they just apply the leader's log in order. This makes single-leader the default choice whenever an application's write volume fits on one machine and doesn't need multi-region write availability. ## The trade-off The trade-offs run in both directions. | Direction | What the topology buys or costs | |---|---| | **Upside** | Read throughput scales roughly linearly by adding more followers, because reads are embarrassingly parallel across replicas. You also get a warm standby for free: if the leader dies, a follower can be promoted, giving you failover-based availability. | | **Downside** | Write throughput is capped by what a single leader machine can handle, since writes cannot be parallelized across nodes without breaking the single-order guarantee. Followers replicate asynchronously in most production setups, which introduces **replication lag** - the gap between when the leader commits a write and when a given follower has applied it. | That lag is the root cause of several read anomalies: a client can write to the leader, then immediately read from a lagging follower and not see its own write; two reads from different followers can even go backwards in time relative to each other if the second read hits a more-lagged replica than the first. ## Failure modes Failure modes in production cluster around three areas. 1. **First, follower lag.** Under write bursts or slow disks, followers can fall minutes behind, and naive load-balancer routing of reads to 'any replica' surfaces stale data to users, often reported as 'I just posted a comment and it disappeared.' 2. **Second, leader failover.** Promoting a follower is not instantaneous - the system needs to detect the leader is actually dead (not just slow, to avoid split-brain with two nodes both thinking they're the leader), pick the most up-to-date follower, and repoint clients, during which writes are unavailable; if the old leader had accepted writes that never replicated before it died, those writes are silently lost unless there's a reconciliation step. 3. **Third, network partitions** between leader and followers can cause replication to stall entirely while the leader keeps accepting writes, ballooning the lag until the partition heals or a follower is promoted and diverges permanently from the old leader. ## Where it shows up A concrete, widely used example is PostgreSQL streaming replication or MySQL's binlog-based replication: one primary handles all `INSERT/UPDATE/DELETE` traffic and ships its `WAL/binlog` to one or more read replicas, which are commonly placed behind a read-only connection string used by reporting jobs and read-heavy API endpoints, while the primary connection string is reserved for transactional writes. MongoDB's replica sets and Kafka's per-partition leader/follower model follow the same pattern at different layers of the stack: - one authoritative writer, - N followers replaying its log, - and automated leader election (via a heartbeat/voting protocol) when the primary disappears. Understanding this topology is foundational because multi-leader and leaderless replication are best understood as answers to the specific limitations - write throughput ceiling and failover unavailability - that single-leader replication accepts as its cost of simplicity.
- If a follower in a single-leader setup falls 30 seconds behind, what specifically can go wrong for a user who just submitted a write?The user's own write may not appear if their next read is routed to that lagging follower, violating read-your-writes; also, if they refresh multiple times and get routed to different followers with different lag, the data can appear to move backwards in time, violating monotonic reads. Applications typically fix this by pinning a user's reads to the leader (or a replica known to be caught up) for a short window after they write.
- How does the system decide which follower to promote when the leader fails, and what can go wrong?Typically the most up-to-date follower (highest applied log position) is chosen via a leader-election protocol with heartbeats and a quorum vote to avoid two nodes believing they're leader simultaneously (split-brain). What can go wrong: if the detection timeout is too short, a merely slow leader gets falsely demoted, causing an unnecessary failover storm; if too long, the system stays write-unavailable longer than necessary.
- Why can't you just scale writes by adding more leaders in this topology?Because the entire point of single-leader is a single total order for writes; adding a second writable node reintroduces the conflicting-write problem this topology exists to avoid. Scaling writes further requires either a bigger single machine, sharding the data across multiple independent single-leader clusters, or switching to multi-leader/leaderless replication with conflict resolution.
Like a single editor-in-chief who approves every change to a shared document and then distributes the approved version to regional offices - the offices can each serve local read requests from their copy, but they're always a little behind the editor's desk, and if the editor's office burns down someone has to be promoted before new edits can be approved again.
saying these in an interview costs you the question
- Says replicas can accept writes in single-leader replication
- Doesn't know followers can lag or that lag causes any user-visible issue
- Assumes reads are always up to date regardless of which node serves them
- Thinks failover is instantaneous with zero write unavailability
- Confuses this with sharding/partitioning