skip to content

questions

6

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?

level: juniorimportance: must knowfreq 75%

answer

  1. one writer, many readers
  2. WAL/binlog streamed to followers
  3. leader dies -> promote a follower
  4. lag causes stale reads

basics

~20 s

One 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 s

Single-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

for a junior

Should describe the basic shape: one leader takes writes, followers get copies, reads can go anywhere. Doesn't need failover mechanics.

for a middle

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.

for a senior

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.

for a principal

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

context

open as a page

What causes replication lag between a leader and its followers, and how would you design a system to guarantee 'read-your-writes' consistency for a user despite that lag?

level: middleimportance: must knowfreq 80%

basics

~20 s

Followers apply changes slightly after the leader because copying and applying takes time, especially under load. Read-your-writes means a user should always see their own recent changes - you guarantee it by making sure that right after someone writes, their next read goes to the leader or to a follower confirmed to be caught up, instead of a random possibly-stale replica.

open as a page

In a single-leader replication setup, what's the practical difference between synchronous and asynchronous follower replication, and what does each cost you?

level: middleimportance: must knowfreq 70%

basics

~20 s

Synchronous means the leader waits for a follower to confirm it got the write before telling the client 'done' - safer but slower. Asynchronous means the leader tells the client 'done' immediately and sends the copy in the background - faster but a crash can lose the last few writes.

open as a page

When would a team choose multi-leader replication over single-leader, and what fundamentally new problem does allowing more than one writable node introduce?

level: seniorimportance: must knowfreq 55%

basics

~20 s

Multi-leader means more than one server can accept writes, usually one per data center or region, so users write to a nearby server instead of one far away, and the system stays writable even if one region goes offline. The catch: two leaders can accept conflicting writes to the same piece of data at nearly the same time, and someone has to decide which one wins.

open as a page

In a leaderless (Dynamo-style) replication system where any of N replicas can accept a write, how do write quorums (W) and read quorums (R) work together to make reads see recent writes, and what do read repair and hinted handoff do?

level: seniorimportance: should knowfreq 45%

basics

~30 s

There's no single leader - a client writes to several replicas at once and only needs a certain number (W) to confirm before it's considered done; reads similarly query several replicas (R) and use the newest answer. If W+R is more than the total number of replicas, at least one replica in any read overlaps with one in the write, so a recent write is very likely seen. Read repair fixes replicas caught with stale data during a read; hinted handoff lets a temporarily unreachable replica's write be held by another node and delivered later.

open as a page

What is the 'monotonic reads' anomaly in a replicated system, what specifically causes it, and how would you prevent a single user's session from experiencing it?

level: principalimportance: should knowfreq 35%

basics

~20 s

Monotonic reads means once you've seen a piece of data, you should never later see an older version of it. The anomaly happens when a user's repeated reads land on different followers with different amounts of lag - a later request can hit a more-behind replica and appear to go 'back in time.' Fix: always route one user's reads to the same replica for their session.

open as a page