Concept lesson · Foundations
Replication and durability
Start here
Definition
Replication maintains copies of the same logical data on multiple machines. Durability is the promise that a successfully committed change survives a stated set of failures. Replication can help provide durability, but the write acknowledgment and recovery rules determine what actually survives.
Why it matters: One machine can fail. Copies can keep data available and spread reads, provided we know which copy is authoritative and when a write is safe to acknowledge.
The leader acknowledges cart version 41 after the protocol commits it on two durable copies. Safe elections must preserve that committed history. A follower may still serve an older applied value.
Read the diagram step by step
- The leader appends version 41 to its durable log; follower B records it durably and acknowledges the leader.
- The leader commits and acknowledges according to a protocol whose election rules preserve committed entries. Counting two copies alone does not prove this property.
- Follower C has applied only version 40. A read there can be stale even while committed version 41 survives one copy loss under the stated protocol.
- Durable bytes, safe failover, and read freshness are separate guarantees.
Worked example
A stores cart version 41 and replies before B receives it. If A is permanently lost, B only has version 40. Waiting for the required durable replica acknowledgments closes that particular loss window, at the cost of latency and write availability.
Key takeaways
- A received update is not necessarily durable or queryable.
- Synchronous and asynchronous replication trade acknowledgment delay against loss exposure.
- Replicas copy bad writes too; backups preserve earlier history.
You will learn to
- Distinguish extra copies from the protocol that updates them.
- Identify exactly which failures an acknowledgment promises to survive.
- Explain stale reads, safe leader replacement, and recovery from a replicated mistake.
Practice in this chapter
8 interview questions with model answers and follow-ups.
Go to interview practiceUseful foundations: Distributed systems: scalability, reliability, availability and efficiency
Workload and timing examples are interview assumptions.
Dotted concept links open the relevant explanation in a new tab.
01What are replication and durability?
Replication means maintaining copies of the same logical data on several machines. A replica is one of those copies; engineers also use the word for the database instance that holds it. Durability means a committed change survives the failures covered by the system's guarantee. A replica can exist and still be too far behind to preserve an acknowledged write.
In single-leader replication, one leader orders writes and followers copy its log. In multi-leader replication, multiple leaders accept writes, so concurrent changes need a conflict rule. In leaderless replication, clients or coordinators contact multiple replicas; versioning, quorums, and repair determine the result. This lesson first traces the single-leader case because it makes it easy to see when the service may safely tell the client that a write succeeded.
A single-copy database can acknowledge a write that later disappears with the only usable storage. For a bounded example, cart C17 changes from version 40 to version 41 with mugs = 2. The acknowledgment policy must specify whether that result survives a process crash, disk loss, or loss of a complete replica; merely adding machines does not establish the promise.
Redundancy means having additional resources: another database copy, application instance, or network path. Replication is the process that carries changes between copies. Adding an empty second database provides neither a current cart nor a useful recovery path. We need a protocol for moving updates and deciding which state is authoritative.
Assume three database participants, A, B, and C, in separate failure zones. A currently orders writes. Our chosen promise is that an acknowledged cart change survives one participant’s failure. All timestamps are illustrative rather than measurements of a product. The version-41 trace tests acknowledgment, read visibility, and recovery separately.
| Replication topology | Write path | Main coordination problem |
|---|---|---|
| Single leader | One authority orders a shard's writes | Replace it safely; followers can lag |
| Multiple leaders | Several authorities accept writes | Resolve concurrent changes; a local success may later conflict |
| Leaderless | A coordinator contacts the selected replicas | Define versions, quorum membership, reconciliation, and repair |
These describe who accepts and orders writes, not a universal consistency level. A last-writer-wins conflict rule may discard one concurrent cart edit; merging a set of product IDs cannot by itself preserve a quantity decrement. Choose conflict semantics from the operation, not merely from the desire to write locally.
02Replication logs: received, durable, applied, committed
A replication log is an ordered sequence of changes that replicas can receive, persist, and replay. Each entry identifies a change and its position in that history. Queryable data pages are a separate representation, so durable logging and visible application need not happen simultaneously. In the example, A records the C17 version-41 update before forwarding the log entry.
| Participant at 10:00:00.008 | Stored log | Queryable cart |
|---|---|---|
| A | v41 is durable | v41 |
| B | v41 is durable | May still show v40 until replay |
| C | Catching up | v40 |
The protocol determines when an entry is committed: accepted into the authoritative history under its safety rules. Counting network receipts without those rules does not establish commitment.
03Synchronous versus asynchronous replication
The acknowledgment policy chooses how much replication must finish before the client hears “saved.” Synchronous replication waits for a configured stage at designated replicas; asynchronous replication allows that work to continue after the reply. In our three-participant example, a majority is two participants, including the leader. Their durable acknowledgments matter only within a protocol that preserves the resulting committed history.
For the one-replica-loss requirement, choose a protocol that commits after the required durable majority acknowledgment. A waits until B confirms durable receipt at .008, then returns version 41. A subsequent failure of A leaves the committed information on B, and the election/recovery rules must preserve it.
Waiting costs remote network and storage time. It can also prevent writes when the required participants are unreachable. Waiting for every replica often worsens tail latency compared with an appropriate majority protocol. Choose the acknowledgment rule from the failure promise, not from a claim that more copies are always better. PostgreSQL’s standby documentation illustrates configurable acknowledgment stages.
The example is a safe majority protocol, not a claim that any database becomes Raft by waiting for a standby. In PostgreSQL, configured synchronous standbys and synchronous_commit=on wait for remote durable logging; remote_write can stop at the standby operating-system buffer, and remote_apply additionally waits for replay. Promotion eligibility and prevention of competing primaries still require a failover design. An acknowledgment setting does not supply that design.
Our three failure zones protect the stated single-participant loss. If all three are within one region, their count does not establish region-loss durability. Cross-region copies add network delay and require a separate placement and acknowledgment decision.
- 1 → 210:00:00.000 update C17Request: add two mugs → A: durable log v41
- 2 → 3replicate and persistA: durable log v41 → B: durable copy v41
- 2 → 4catch upA: durable log v41 → C: catching up
- 3 → 210:00:00.008 durable acknowledgmentB: durable copy v41 → A: durable log v41
- 2 → 5reply after commitA: durable log v41 → Response: saved v41
- 2 → 6separate retained historyA: durable log v41 → Retained cart history
04Replica lag and read-your-writes
Replica lag is the gap between a source’s progress and a follower’s received or applied state. A read can therefore be stale even after the write commits durably. In the example, a read at .010 seconds reaches C, which still serves v40 although v41 has committed elsewhere. Durability and read visibility require separate policies.
Read-your-writes requires a routing or version mechanism when replicas lag; replication by itself is insufficient.
Remember: Carry the required version; do not silently return an older one.
A commit-position token identifies the write’s location in a particular replication history. A follower’s applied position identifies how far it has replayed that same history into queryable state. Comparing those positions lets a read wait for its required write instead of guessing how many milliseconds replication needs.
After a write, send the read to the current leader. Alternatively, return the write’s log position and make a follower wait until it has applied that position before answering. Use the database’s supported mechanism; an application-assigned version number alone cannot prove a follower has caught up.
| Read policy | Benefit | Cost |
|---|---|---|
| Read the authority | Straightforward current-state path | Concentrates reads and needs reachable authority |
| Wait for a verified position | Distributes session reads | Waiting and token handling |
| Read any follower | Low local latency | An explicit stale-read contract |
For the cart, we choose read-your-writes: the client should observe its own acknowledged change. Product browsing can use a different policy. A fixed sleep is only a guess because lag can grow under load or failure.
A current-state read needs more than a machine that once accepted writes. A linearizable read must fit an order that respects completed operations in real time, so it cannot return a version from before a write that completed before the read began. Checking current authority is part of establishing that guarantee after failover.
A server calling itself “leader” may be an isolated former leader with old data. For linearizable reads, a consensus-based database must confirm its current leadership and apply the required committed entries before answering. A read-your-writes token must also refer to the correct history after failover. A number from another shard or a discarded history does not prove this replica includes the write.
05Leader failover and split-brain prevention
Leader failover lets another replica accept writes when the leader becomes unusable. Missing replies cannot tell us whether A crashed or lost its network connection. If B takes over while A keeps accepting independent writes, their data can diverge: this is split brain. The election and storage protocol must prevent it. Consider A becoming unreachable after the version-41 commit.
Safe failover has four obligations:
In this example, B and C form the required majority and the protocol establishes B as leader. A’s later return does not automatically restore its authority.
Response loss requires operation deduplication independently of replication. If v41 committed but its reply was lost, retry with the same operation identifier and recover its durable outcome. Applying “add two mugs” twice would turn uncertainty into four mugs. Routing clients away from A changes discovery; it does not itself fence obsolete writes.
06Replication versus sharding versus backups
Three copies of C17 are replicas. Three servers each holding different customers are shards. Replication helps survive loss and may add read capacity; sharding divides data and work. Adding followers does not automatically multiply a single leader’s write capacity, because every follower still processes the write stream.
Letters stand for records. Compare the copies, subsets and timestamps instead of memorizing three labels.
Remember: Copy the present, split the present, or retain the past.
Try from memoryIf a bad delete reaches every replica, which arrangement can restore the old record?
A suitable retained backup or recovery history. Replication alone can faithfully copy the bad delete.
Copies must occupy appropriate failure domains. Three processes on one laptop do not survive laptop loss. Three zones still share risks such as a bad application release or administrator action.
07Replica repair: anti-entropy, Merkle trees, read repair, and hinted handoff
Replica repair detects and reconciles differences between copies that missed updates. The repair must follow the store's version and conflict rules; it cannot simply trust whichever machine responds first. A leader/follower log normally catches up by replaying missing committed entries or installing a snapshot. The following mechanisms are common in Dynamo-style replicated stores and must not be confused with electing a new leader.
| Mechanism | Trigger and action | Main limit or cost |
|---|---|---|
| Hinted handoff | A coordinator retains a missed update for an unavailable replica and replays it later | Temporary, best-effort delivery; hints can expire or their holder can fail |
| Read repair | A read encounters different versions and repairs the data participating in that read | Cold data may never be read; blocking repair adds read latency |
| Anti-entropy repair | A background or scheduled process compares replicas over shared ranges and transfers differences | Scans and streaming consume disk/network capacity; complete coverage must be verified |
| Merkle tree | Hierarchical hash summaries locate differing subranges | Detects differences; it does not choose the correct version or resolve a business conflict |
Why repair must cover cold data
Anti-entropy means systematically reducing divergence, including records that receive no foreground reads. Suppose A and B contain cart C17 at version 41 while C still has version 40. A surviving hint may deliver the missed update to C. A read comparing B and C may repair that particular cart. Scheduled range repair also discovers the difference when nobody reads C17. Hints and read repair therefore reduce inconsistency but do not replace full repair coverage.
How a Merkle tree locates differences
To compare replicas without first transferring every record, compute compact hash summaries of the same data ranges. A hash is derived from encoded bytes, so the replicas need a canonical encoding: the same record must produce the same byte representation on both machines. The tree then organizes those summaries so a mismatch can be narrowed to a smaller range.
A Merkle tree summarizes data from the bottom up: leaves hash canonical records or small ranges, and each parent hashes its children. Compare roots for the same range and comparable repair snapshot. If they differ, descend only into mismatching branches. For four leaf ranges, matching left-half summaries let replicas focus on the right half containing C17 instead of transferring every record. After locating differences, exchange the actual versioned data and reconcile it. Matching hashes are equality evidence under the chosen collision assumptions, not a mathematical guarantee of uniqueness. Building the summaries still costs work, even when little data needs streaming.
Why deletion evidence must survive
Deletes require repair too. A tombstone is a versioned deletion marker that tells another replica its older value must remain deleted. If A and B delete C17 while C is offline, immediately erasing both the value and its tombstone removes that evidence. When C returns with version 40, repair could resurrect the deleted cart. Retain deletion evidence long enough for every relevant replica to be repaired, or exclude and rebuild a replica that missed the safe recovery horizon. In Cassandra, plan and verify repair completion before the applicable gc_grace_seconds horizon; actual tombstone removal also depends on compaction and table settings. Time passing alone does not prove that every replica learned the delete.
Repair promotes convergence under its delivery, retention, and conflict-resolution assumptions. It does not undo a stale response already returned, recover a write absent from every surviving copy, or establish linearizability by itself. Cassandra's blocking read repair supports a specific monotonic-quorum-read behavior; it is not a general transaction guarantee. Monitor completed range coverage, repair age, hint backlog, and repair resource use instead of treating a started repair job as proof of recovery.
08Interview answer: defend the acknowledgment and read policy
Interviewer: “Why not just put a read replica behind the load balancer?”
Candidate: “That can improve reads, but I first need to define what an acknowledged cart update survives. With asynchronous replication, A could acknowledge v41 and fail before B receives it. For our one-node-loss promise, I would choose the necessary durable acknowledgments and a safe election protocol.
“I would also avoid sending a read-after-write request to an arbitrary lagging follower. An authoritative read or verified replication position preserves the session’s read-your-writes contract. Finally, if a buggy job deletes the cart, every live replica may copy that deletion. I need retained history and a restore procedure for that failure.”
This answer separates keeping an acknowledged change, showing the right version, and recovering an earlier valid state. It explains why the extra components exist and what each costs. A diagram of repeated databases becomes useful once those behaviors are explicit.
Practise the interview questions
Say your answer aloud before opening the model answer. Then answer the follow-up and compare the reasoning.
What is replication? How is it different from redundancy and durability?
Reveal a model answer
Replication copies changes to additional replicas. Redundancy is the broader idea of spare resources: a spare machine, disk, or network link can be redundant without containing a usable data copy. Durability is the guarantee that a committed write survives a defined failure set. The replication protocol, durable storage, acknowledgment rule, and failover rules jointly determine that guarantee.
Suppose leader A acknowledges cart v41 before follower B receives it. Replication is configured, but permanently losing A can still lose that acknowledged write. Waiting for the required durable copies reduces this loss exposure while adding network/storage latency and making writes depend on those copies being reachable. Replication also copies a mistaken deletion, so it does not replace a backup.
Interviewer follow-up
Could an old backup count as redundancy?
Reveal the follow-up answer
“It is another copy, but a day-old backup supports a different promise from preserving a cart change acknowledged this minute.”
What the answer must demonstrate: Name the freshness and failure promise.
Why distinguish received, durable, and applied?
Reveal a model answer
“Received bytes may be only in memory. Durable bytes survive the specified storage failure model. Applied entries are visible to queries. B can have v41 durably logged while ordinary reads still show v40, so acknowledgment and read policy must account for different milestones.”
Interviewer follow-up
Should a save wait until every replica applies?
Reveal the follow-up answer
“Only if that is the required contract. I can use a suitable durable commit rule and separately route or wait for reads that must see the write.”
What the answer must demonstrate: Do not equate a network acknowledgment with query visibility.
A leader acknowledges v41 before a follower receives it, then permanently fails. Explain the possible data loss.
Reveal a model answer
“A responds at .003, fails at .006, and B would receive the change at .008. If A’s storage is lost, the survivors have v40. I either accept that acknowledged-write loss window explicitly or wait for the required durable replica before answering.”
Interviewer follow-up
Does synchronous replication prevent every kind of loss?
Reveal the follow-up answer
“No. It covers a stated failure model. Correlated storage loss, a replicated bad delete, or unsafe recovery can exceed it.”
What the answer must demonstrate: Avoid universal durability claims.
With three replicas and a one-replica-loss durability goal, why might the commit protocol wait for two durable copies instead of all three?
Reveal a model answer
“Two durable copies leave at least one copy of an acknowledged entry after any one participant is lost. With a safe election and commit protocol, the surviving majority preserves that committed history and can continue. Waiting for all three adds a copy but makes the slowest replica control acknowledgment and stops writes if any replica is unreachable. I would choose two only because it meets the stated one-failure contract; the count alone is not the safety proof.”
Interviewer follow-up
What if two simultaneous storage losses must be tolerated?
Reveal the follow-up answer
“I must revisit replica count, acknowledgment, and placement together. One surviving copy cannot preserve a write it never received.”
What the answer must demonstrate: Failure budget and acknowledgment must agree.
A write of v41 succeeds, but a subsequent session read returns v40. What should you inspect?
Reveal a model answer
“Check which replica answered and how far it had applied the write log. It may have saved v41 without making it readable yet. To read my own write, use the verified current leader or wait for a follower to apply the returned commit position. That position must still identify the right history after failover. A former leader or an arbitrary application version cannot prove freshness.”
Interviewer follow-up
Why is 100 ms of waiting insufficient?
Reveal the follow-up answer
“Lag is not bounded by that guess during overload or failure. I need evidence that the required update became visible.”
What the answer must demonstrate: Waiting a fixed time does not prove that the required update is visible.
How do replicas and shards fit together?
Reveal a model answer
“A shard owns a subset of records, while replicas store copies of that subset. Cart C17 can belong to one shard with three replicas. Adding shards can divide data and write work; adding followers preserves copies and can spread eligible reads. Each follower still has to process its shard’s write stream.”
Interviewer follow-up
Will ten followers give ten times the write capacity?
Reveal the follow-up answer
“Not by themselves. A single leader still orders the stream, and each follower must keep up with it.”
What the answer must demonstrate: Do not count duplicated processing as partitioned work.
A new leader B takes over from isolated leader A. What prevents A from continuing to commit writes?
Reveal a model answer
“A must lose the ability to commit new writes when B takes over. Missing heartbeats alone does not prove A stopped. Use the database’s safe election and fencing protocol to reject the old leader, then update routing so clients find B.”
Interviewer follow-up
Reveal the follow-up answer
“No. Cached addresses and existing connections may still reach A. The protected write path must reject obsolete authority.”
What the answer must demonstrate: Routing is discovery, not ownership enforcement.
Every replica contains a mistaken deletion. What next?
Reveal a model answer
“I stop the faulty job, restore retained history in isolation, identify C17’s last valid state, and verify the repair. Promoting another current replica cannot undo a deletion they all copied correctly. I would also check the full affected range.”
Interviewer follow-up
What proves that recovery plan works?
Reveal the follow-up answer
“A measured restore that validates records and application behavior within the objectives. Backup completion alone proves only part of the path.”
What the answer must demonstrate: Replicas and recovery history solve different failures.
Blank-page exercise · 15 minutes
Build the answer yourself
Specify received, durable, applied, and committed milestones for a three-replica protocol. Test cart C17 version 41 by failing leader A before and after acknowledgment, then evaluate a stale read and a lost-response retry.
- Label received, durable, applied, and committed separately.
- Show which surviving participant contains v41.
- Explain a read-after-write request and an ambiguous retry.
- Demonstrate why a replicated bad deletion needs retained history.
Check that each component and design decision follows from your requirements and workload.
Recall the key ideas
Answer from memory before opening each card. Explain why the choice works and what it costs. Revisit missed cards tomorrow.
Replication and durabilityWhat must “cart saved” mean?Recall first, then reveal
State where the write must be durably stored before success is returned and which failures it must survive. For example, a safe commit and election protocol can require durable storage on two of three replicas to tolerate one replica loss.
Saved where, before saying saved.
Return to lessonReplication and durabilityWhy might a durable follower show v40?Recall first, then reveal
It may have stored v41 in its log without applying it to queryable state yet.
Received, stored, applied: three milestones.
Return to lessonReplication and durabilityAre replicas a replacement for backups?Recall first, then reveal
No. A bad delete can reach every live copy; retained history is needed to recover the earlier state.
Replicas copy today; backups preserve yesterday.
Return to lessonReplication and durabilityDoes adding read replicas divide the write stream?Recall first, then reveal
No. Copies still process the same writes; sharding divides different data and work.
Copy versus divide.
Return to lessonFinal revision
Summary and interview notes
Replication keeps copies; durability defines which committed changes survive which failures. The acknowledgment rule, safe failover protocol, read policy, and retained recovery history must be chosen together.
Remember these points
- Received, durable, applied, and committed are different milestones.
- Asynchronous replication can lose an acknowledged write if its only durable copy is lost before followers catch up.
- A safe two-of-three majority protocol can preserve committed history through one participant loss; copy counting alone cannot.
- Define how old follower reads may be. Before trusting a leader’s read, verify it is still the leader.
- Replicas help recover from component loss; retained backups and logs help recover from replicated mistakes.
Interview tips
- For every successful write, point to the surviving durable copy after the failure you claim to tolerate.
- Test a lost response, a lagging read, and an isolated former leader separately; each needs a different mechanism.
- Name failure domains explicitly: process, disk, zone, and region are not interchangeable.
Important qualifications
- Synchronous replication settings do not automatically select a safe replacement or fence the former primary.
- A read-your-writes token must identify the relevant committed update even after failover; a sequence number from an unrelated or discarded history is insufficient.
Technical references
- PostgreSQL: Log-Shipping Standby ServersVerified reference for asynchronous streaming and distinct remote acknowledgment stages.
- Raft: In Search of an Understandable Consensus AlgorithmPrimary reference for safe leader election and preserving committed history. The cart trace is hypothetical.
- Apache Cassandra: RepairOfficial range repair, Merkle summaries, repair coverage, and repair-before-tombstone-expiry guidance. Checked 2026-09-23; the page identifies its documentation version as 5.0.
- Apache Cassandra: HintsOfficial explanation of coordinator hints, later handoff, and why best-effort hints do not replace anti-entropy repair.
- Apache Cassandra: Read RepairOfficial read-repair scope and blocking/none tradeoffs; monotonic quorum reads are narrower than general linearizability or transaction isolation.
- Apache Cassandra: TombstonesOfficial deletion-marker, resurrection, grace-period, and compaction-removal conditions; no automatic-safe-GC claim.
- Dynamo: Amazon’s Highly Available Key-value StorePrimary section 4.7 explains Merkle-tree anti-entropy and hinted-handoff limitations; these are distinct from safe leader election.
Practice marks stay in this browser.