System designby Learnastra

Concept lesson · Foundations

Quorums, consensus, leases, and fencing

By Anup Rai

Start here

Definition

A quorum is a protocol-defined set of participants whose votes or replies are sufficient for an operation to proceed, often a majority. Consensus is a protocol for agreeing on a value or ordered history despite specified failures. A lease grants time-limited authority; fencing makes the protected resource reject obsolete authority.

Why it matters: A replacement leader or worker must be able to take over without letting an isolated or paused old owner corrupt the result. Counting responses, agreeing on ownership, and enforcing ownership are separate jobs.

The visual modelQuorum intersection and stale-writer fencing

A majority intersects every other majority. The resource must still reject a stale worker token.

Quorum intersection and stale-writer fencingA majority intersects every other majority. The resource must still reject a stale worker token. With three voters, majorities {A,B} and {B,C} share B. Consensus uses rules beyond this overlap to agree on a log. Worker W1 once held fencing token 7. W2 takes over with token 8. The protected store remembers 8 and rejects W1 with token 7 even if W1 wakes after its lease expired. An expired lease alone cannot stop code already running on a paused machine.Agree on an owner; reject the obsolete writerABCmajority {A,B}majority {B,C}Shared voter BOverlap is necessary;consensus adds log rules.W1 resumestoken 7W2 owns job E9token 8Store: max token 8accept token 8reject token 7write 8write 7: rejectedThe store must check the token atomically with the write.
Read the diagram step by step
  1. With three voters, majorities {A,B} and {B,C} share B. Consensus uses rules beyond this overlap to agree on a log.
  2. Worker W1 once held fencing token 7. W2 takes over with token 8.
  3. The protected store remembers 8 and rejects W1 with token 7 even if W1 wakes after its lease expired.
  4. An expired lease alone cannot stop code already running on a paused machine.

Worked example

Of three controllers, two agree to replace worker W1 (epoch 7) with W2 (epoch 8). W2 publishes with token 8. When W1 resumes and presents token 7, the output store rejects the stale write.

Key takeaways

  • R + W > N proves read/write set overlap, not linearizability by itself.
  • Consensus establishes committed authority; a lease expiring cannot stop a paused process from resuming.
  • A fencing check must be atomic with the protected write at the resource.

You will learn to

  • Calculate quorum overlap and explain what it does not prove.
  • Describe how an agreed log preserves one ownership history.
  • Show why a resource must reject obsolete ownership even after a lease expires.

Practice in this chapter

8 interview questions with model answers and follow-ups.

Go to interview practice

Useful foundations: Replication and durability · CAP theorem: consistency, availability, and partition tolerance

Workload and timing examples are interview assumptions.

01What are quorum, consensus, lease, and fencing?

A service that replaces failed workers has two problems: agree which worker is now authorized, and prevent a previous worker from publishing afterward. The following mechanisms handle different parts of that handover.

Keep four definitions separate:

  1. Quorum: a protocol-defined set of participants whose votes or replies are sufficient for an operation to proceed, often a majority.
  2. Consensus: agreement on a decision or ordered history despite failures within the protocol’s model.
  3. Lease: permission that expires after a defined time.
  4. Fencing token: a monotonically increasing ownership number checked by the protected resource to reject outdated writers.

These are not four names for a distributed lock. Quorum rules can require read and write groups to share a replica. Consensus makes the controllers agree on the sequence of ownership changes. A lease limits permission in time. Fencing lets the output store reject a write carrying an older ownership number. Start by keeping those responsibilities separate.

Publishing the result of a job shows why agreeing on its owner and protecting its output are separate tasks. Export E9 reads records, writes an output file, and publishes its location in a manifest. One worker can own the job initially, but recoverable execution needs a replacement owner after failure. The replacement decision and rejection of stale publication are separate requirements.

We introduce an ownership record: E9 → worker W1, epoch 7. An epoch is a number that increases whenever the job is assigned a new owner. W1 may work while its time-limited permission, called a lease, remains valid. Three controller replicas store this permission so losing one controller need not lose the job’s ownership history.

At 12:00:02 W1 pauses for twelve seconds. The controllers later expire its permission and assign W2. A paused process is not dead: W1 can resume with its old instructions. Our design must answer two different questions: how do controllers agree on the new owner, and how does the output store stop the old owner? All times and numbers in this lesson are hypothetical.

02Quorum arithmetic: N, R, W, and overlapping sets

Concept in focusQuorum overlap is set intersection

With N = 5 and R = W = 3, every such read set intersects every such write set. These sets illustrate arithmetic, not a complete consensus algorithm.

Quorum overlap is set intersectionWith N = 5 and R = W = 3, every such read set intersects every such write set. These sets illustrate arithmetic, not a complete consensus algorithm. There are five replicas A through E. A completed write uses A, B and C; a read uses C, D and E. The shared C illustrates why any read and write sets overlap when R + W > N. Version selection and concurrency rules are still needed for a consistency guarantee.N = 5 replicas; W = 3 writes; R = 3 readsABCDEWrite set {A, B, C}Read set {C, D, E}Overlap at C: R + W > N forces any such read set to intersectthe completed write set.Intersection alone is insufficient. The protocol must still resolve versions,concurrent writes and failures correctly.

Remember: Overlap finds a shared participant; the protocol makes its evidence useful.

Read the diagram
  1. There are five replicas A through E.
  2. A completed write uses A, B and C; a read uses C, D and E.
  3. The shared C illustrates why any read and write sets overlap when R + W > N.
  4. Version selection and concurrency rules are still needed for a consistency guarantee.

For N = 3, W = 2, and R = 2:

Write group holding version 8 Possible read group Shared participant
R1, R2 R1, R2 R1 and R2
R1, R2 R1, R3 R1
R1, R2 R2, R3 R2

Every read has a chance to encounter the acknowledged version. If W is also greater than N/2, any two write groups overlap. One unavailable controller still leaves two participants, but two unavailable controllers leave too few for these operations. These counts describe the chosen fixed membership and response requirements.

There are two counts to keep separate. For simple majority consensus, N = 2f + 1 participants can continue with f unavailable when the remaining majority communicates and the protocol's timing assumptions eventually hold. Four voters still need three votes and tolerate only one unavailable voter; five need three and tolerate two. Membership changes must themselves follow the protocol: changing N independently on different clients invalidates the fixed-set intersection argument.

03Why quorum overlap alone is not a consistency protocol

Suppose a failed update proposing owner version 9 reaches only R1. One read consults R1 and R2 and completes with 9. Only after that response, another read starts, consults R2 and R3, and returns 8; no new ownership update occurred between these reads. Both read groups have size two, yet clients have observed a reversal unless the protocol handles that incomplete write correctly.

We need rules for valid versions, concurrent updates, failed attempts, and read completion. “Take the largest timestamp” is not automatically correct: clocks can disagree and an incomplete proposal may not be committed. Using substitute nodes during a failure also changes the overlap assumptions. This is one reason Dynamo’s quorum-style techniques must be understood with their surrounding protocol.

For E9’s ownership, we want one agreed committed history. We therefore choose an established consensus protocol rather than invent a lock service from the arithmetic alone. A quorum contributes to the proof; it is not the whole proof. The extra discipline costs coordination and can stop progress without enough connected participants.

A pending write is allowed to take effect even if its caller never receives success. The error in the trace is returning 9 and then reverting to 8 with no intervening write. Some atomic read/write-register protocols address this by making a reader propagate the selected version to a quorum before returning. The Attiya–Bar-Noy–Dolev register is a classic example. That is a different protocol from a one-round 'read two and return the maximum' rule, and from consensus on arbitrary ownership commands.

A sloppy quorum may acknowledge on substitute nodes outside a key’s normal replica set when home replicas are unavailable. For home replicas A/B/C, two substitutes D/E can accept a write while a read of A/B sees the old value. Counting W=2 and R=2 against N=3 does not prove overlap because those responses came from different sets. Hinted handoff can later deliver the missed data to home replicas. This improves write availability under the chosen contract, but adds repair work and does not supply an immediate latest-value read guarantee.

Worked example diagramController replicas agree that W2 owns export E9 at epoch 8. The output store accepts W2’s epoch-8 publication and rejects W1’s delayed epoch-7 write.
Quorums, consensus, leases, and fencing: architecture diagram1. Controllers R1/R2/R3 to 2. W1: old epoch 7: earlier ownership expires; 1. Controllers R1/R2/R3 to 3. W2: new epoch 8: 12:00:11 grant epoch 8; 3. W2: new epoch 8 to 4. Output store: latest fence 8: 12:00:12 publish with fence 8; 2. W1: old epoch 7 to 4. Output store: latest fence 8: 12:00:14 stale fence 7 rejected; 4. Output store: latest fence 8 to 5. Published E9: retain accepted output1 → 2: earlier ownership expires1 → 3: 12:00:11 grant epoch 83 → 4: 12:00:12 publish with fence 812:00:14 stale fence 7 rejected4 → 5: retain accepted output01Controllers R1/R2/R302W1: old epoch 703W2: new epoch 804Output store: latestfence 805Published E9
  1. 1 → 2earlier ownership expiresControllers R1/R2/R3 → W1: old epoch 7
  2. 1 → 312:00:11 grant epoch 8Controllers R1/R2/R3 → W2: new epoch 8
  3. 3 → 412:00:12 publish with fence 8W2: new epoch 8 → Output store: latest fence 8
  4. 2 → 412:00:14 stale fence 7 rejectedW1: old epoch 7 → Output store: latest fence 8
  5. 4 → 5retain accepted outputOutput store: latest fence 8 → Published E9

04Consensus with Raft: leaders, terms, and committed logs

Consensus lets a group agree on state transitions under a defined failure model. In a replicated-log approach, replicas apply the same committed commands in the same order. For E9, that sequence includes assigning W1, expiring its ownership according to the lease policy, and granting W2 epoch 8.

Raft organizes this around a leader, followers, and election terms. A term is a generation of controller leadership, distinct from E9’s job-ownership epoch. The leader replicates log entries; election and commit rules preserve committed history across leader changes. A majority of three is two; a majority of five is three. These are crash-fault protocols, not a claim that any malicious participant can be tolerated. Raft paper.

At 12:00:11, the controllers agree that W2 owns the job under epoch 8. A controller with old data must not grant the job again. Before reporting the current owner, it must also perform the protocol’s check that its answer is current. A recent timestamp alone cannot prove that.

Raft: election, replication, commitment, application

  1. Elect. A candidate gains a majority with the required log-freshness check.
  2. Replicate. The elected leader replicates an entry from its current term.
  3. Commit. That entry is committed after the specified majority accepts it.
  4. Apply. Replicas apply committed entries in order.

A majority containing an old-term entry alone is not enough to infer that entry is committed. For a linearizable read without appending each read, the leader must establish current authority, know the committed position, and apply through it before answering. These rules explain why the label “leader” is insufficient.

Basic Paxos: prepare, select, accept

Paxos is another consensus protocol. Basic Paxos chooses one value using proposers and acceptors:

A proposer asks the group to choose a value. Acceptors retain promises and accepted proposals so later attempts can discover earlier decisions. A ballot is a uniquely ordered proposal-attempt identifier; a higher ballot gives an attempt priority, not permission to replace an already chosen value.

  1. Prepare. A proposer with a unique higher ballot asks a majority to promise not to accept lower ballots. Replies report previously accepted values.
  2. Select the safe value. Carry forward the value from the highest accepted ballot learned, if any; otherwise propose a new value.
  3. Accept. A majority accepting that ballot/value makes the value chosen. Durable promises and accepted state protect recovery.

The rule for carrying an earlier value forward, together with intersecting majorities, prevents two different chosen values. Repeated competing proposals can prevent progress; practical systems use leadership and sufficient communication stability. Multi-Paxos builds an ordered log from repeated decisions, often amortizing preparation under stable leadership. Like Raft, it is more than majority arithmetic and is distinct from two-phase commit across independent databases.

Interview check: Can a new proposer ignore a previously accepted value because it has a larger ballot? No; the prepare replies constrain the value it may safely propose.

05Fencing tokens: reject stale writers at the resource

The controllers agree that W2 owns E9, but the output store still receives worker requests independently. Each publication therefore includes the agreed ownership version (epoch) as a fencing token. When updating the manifest, the store checks that version atomically. Otherwise, a paused old worker could resume and overwrite W2’s result despite the controllers’ agreement.

W2 finishes quickly. At 12:00:12 it asks the output store to publish file E9-v8 with fence 8. The store atomically compares the fence with its latest accepted ownership generation and records the new manifest. At 12:00:14, W1 resumes and submits E9-v7 with fence 7. The store rejects it because 7 is older than the accepted 8.

Concept in focusFencing rejects the paused old owner

A lease can expire while a process is paused. The protected resource must enforce the fencing token; issuing tokens alone is insufficient.

Fencing rejects the paused old ownerA lease can expire while a process is paused. The protected resource must enforce the fencing token; issuing tokens alone is insufficient. Old worker to Old worker: Worker pauses while holding token 41. New worker to Resource: New owner writes with token 42; resource records the newer token. Old worker to Resource: Old worker resumes and writes with 41. Resource to Old worker: Reject the stale token at the resource boundary.Old workerResourceNew workerWorker pauses while holding token 41.New owner writes with token 42; resource records the newertoken.Old worker resumes and writes with 41.Reject the stale token at the resource boundary.

Remember: The store rejects the older ownership token.

Read the diagram
  1. Old worker to Old worker: Worker pauses while holding token 41.
  2. New worker to Resource: New owner writes with token 42; resource records the newer token.
  3. Old worker to Resource: Old worker resumes and writes with 41.
  4. Resource to Old worker: Reject the stale token at the resource boundary.
Time Attempt Store decision
12:00:11 Controllers grant W2 epoch 8 Ownership history advances
12:00:12 W2 publishes with fence 8 Accept; remember 8
12:00:14 W1 publishes with fence 7 Reject obsolete owner

Persist the highest accepted fence with the manifest so a resource restart cannot forget token 8. Scope that number to the protected job/resource, and accept only tokens issued through an authenticated ownership path; an arbitrary client-supplied large integer is not authority. A repeated token 8 may be valid, so deduplicate its operation separately. This design prevents a lower-generation write after the store has accepted the higher generation; it does not promise that only one worker ever computed an output.

06Leases, fencing, and idempotency solve different failures

An idempotency key solves a different problem: repeating the same valid publication attempt. Fence 8 can be valid for several W2 requests; it does not identify which repeated request is the same operation. Use a stable publication identifier and a conditional manifest change when duplicate effects matter.

Some stores support checking an ownership key in the same transaction as the data update. etcd’s concurrency API exposes ownership keys for that pattern. An unrelated external service is not automatically inside that transaction. Keep consensus on small critical ownership metadata where useful; copying the export’s large file bytes need not pass through the controller log.

Mechanism Protects Does not establish
Quorum intersection Required response sets share a participant Which observed proposal is safe to return
Consensus log One committed sequence of ownership decisions Atomic completion at an unrelated output service
Lease A time-bounded permission under the lease service's rules That a paused process stopped executing
Resource fence Reject an older generation after a newer one is installed Deduplication or instantaneous global revocation
Idempotency key Recognize the same logical operation again Whether its caller still has permission

For time-based permission, state which service evaluates expiry and which clock assumptions the implementation uses. A worker's cached wall-clock check is not a resource-side authorization check. Clock jumps and long pauses are reasons to use a proven lease implementation and have the output store check permission as part of the publication itself.

07Interview answer: a paused worker returns after takeover

Interviewer: “The lease expired, so why can’t W2 just continue?”

Candidate: “The controllers can agree that W2 owns the job while W1 is only paused. When W1 resumes, it may still try to publish its old result. I attach epoch 8 to W2’s request and make the output store check it atomically when saving. The store rejects older epochs. The controllers choose the owner; fencing makes the store enforce that choice.

“If the controllers lose their majority, I would stop granting new ownership under this protocol rather than invent two histories. Existing work must obey its remaining permission and publication rules. After recovery, I would reconcile the committed ownership record, the latest published manifest, and any abandoned files.”

This answer names what the quorum, consensus log, lease, fence, and idempotency key each contribute. None is a general replacement for the others. The failure drill includes controller loss, an isolated controller, a paused worker, and a publication whose response is lost.

Practise the interview questions

Say your answer aloud before opening the model answer. Then answer the follow-up and compare the reasoning.

Foundation · Question 1

Define quorum, consensus, lease, and fencing. Which problem does each solve?

Reveal a model answer

A quorum is a required response set, such as two of three controllers. Consensus makes those controllers agree on an ownership decision or committed log despite the failures it tolerates. A lease gives ownership for a limited interval. Fencing adds an increasing ownership token that the output store checks atomically with a write.

For export E9, controllers grant W1 epoch 7. W1 pauses; the lease expires; controllers agree to grant W2 epoch 8. W2 publishes with token 8. If W1 resumes and presents 7, the store rejects it. The lease did not stop W1's CPU from executing; the fencing check stops its stale effect after newer authority reaches the resource. Quorum overlap helps the agreement proof but does not supply a complete consensus protocol.

What the answer must demonstrate: Start from actual sets rather than a memorized equation.

Applied · Question 2

Why does overlap not prove linearizability?

Reveal a model answer

“With three replicas, read groups {R1,R2} and {R2,R3} do overlap at R2. But suppose only R1 saw an incomplete write of v9 while R2 and R3 still have v8. A first read returns v9 from R1, then a later read through R2 and R3 returns v8 without another write. The problem is that the read exposed a value without preserving it for later reads. Quorum intersection alone does not define safe version selection, write-back, commitment, or recovery.”

What the answer must demonstrate: Use overlapping replica sets and non-overlapping-in-time reads; explain why a selected value must remain visible to later reads.

Foundation · Question 3

What does consensus provide when three controllers assign one owner for export job E9?

Reveal a model answer

“It gives the controllers one agreed sequence of ownership transitions, so W1 expiry and W2’s epoch-8 grant are not independently invented on different copies. I would use a proven replicated-log protocol whose election and commit rules preserve the history after controller failure.”

What the answer must demonstrate: Agreement on metadata does not atomically include every external effect.

Applied · Question 4

What happens when two of three controllers are unreachable?

Reveal a model answer

“Only one remains, so the majority protocol cannot safely advance ownership. I would stop new grants and report reduced availability. I would not let the isolated replica infer that its stale state is now authoritative because it is the only one this client can reach.”

What the answer must demonstrate: Distinguish safety from continued progress.

Applied · Question 5

W1 has fencing token 7; replacement W2 publishes with token 8. W1 resumes. What must the output store check?

Reveal a model answer

“W2’s publication has fence 8, so the output store has atomically recorded that generation with the manifest. W1 arrives carrying 7. The store rejects 7 before changing the protected state, preventing W1 from replacing W2’s newer result.”

What the answer must demonstrate: A separate preflight check leaves a race.

Follow-up · Question 6

Does a fence instantly revoke old work everywhere?

Reveal a model answer

“Not necessarily. A resource comparing against its latest accepted fence learns about generation 8 when that newer authority reaches it. It prevents older writes after that point. If the requirement is immediate revocation everywhere, I need current-ownership validation or another stronger coordinated boundary.”

What the answer must demonstrate: Describe the precise fencing guarantee rather than implying physical process termination.

Foundation · Question 7

A valid worker retries publication with the same fencing epoch 8. Why is an operation idempotency key still needed?

Reveal a model answer

“Epoch 8 says W2 is an eligible owner. It does not distinguish one publication attempt from a retransmission of the same attempt. I use a stable publication ID so a lost response does not create duplicate effects while that ownership is still valid.”

What the answer must demonstrate: Operation identity and authorization are independent checks.

Follow-up · Question 8

What do you reconcile after the controller outage ends?

Reveal a model answer

“I inspect the committed ownership history, the output store’s accepted fence and manifest, and any unfinished files. A worker saying it finished is weaker than the protected publication record. I then resume or retry with stable IDs and valid ownership rather than blindly rerunning every reported job.”

What the answer must demonstrate: Recovery must consult the state that actually governs the external result.

Blank-page exercise · 18 minutes

Build the answer yourself

Act out export E9 with three controller replicas and two workers. Pause the old worker, transfer ownership, then let both attempt publication.

  • Compute which read/write replica sets intersect.
  • Distinguish the controller log’s term from the job’s ownership epoch.
  • Show the output store atomically rejecting fence 7 after accepting fence 8.
  • Explain why an idempotency key is still needed for repeated valid operations.

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.

Quorums, consensus, leases, and fencingWhat does R + W > N establish?Recall first, then reveal

For a fixed set of N replicas, a read contacting R replicas and a write acknowledged by W replicas must share at least one replica when R + W > N. Rules for versions, incomplete writes, and failures are still needed.

Overlap is a building block.

Return to lesson
Quorums, consensus, leases, and fencingWhy does W1 need fencing after its lease expires?Recall first, then reveal

It may resume and execute old instructions; the output store must reject its obsolete ownership.

The resource checks the ticket number.

Return to lesson
Quorums, consensus, leases, and fencingAre Raft term 4 and E9 epoch 8 the same number?Recall first, then reveal

No. One identifies controller leadership; the other identifies ownership of this job.

Different authorities, different generations.

Return to lesson
Quorums, consensus, leases, and fencingDoes a fencing token deduplicate a valid request?Recall first, then reveal

No. A stable operation ID is still needed to identify repeated effects within the same valid ownership epoch.

Who may write is not which attempt this is.

Return to lesson

Final revision

Summary and interview notes

Quorum rules can require response groups to overlap. Consensus commits an agreed history, leases limit permission in time, and fencing makes the output store reject obsolete ownership numbers. None of these alone makes an unrelated external effect exactly once or stops a paused worker from running.

Remember these points

  • R + W > N assumes the same fixed replica membership; it does not define safe read selection or failed-write handling.
  • A 2f + 1 majority group tolerates f unavailable participants for progress only when the surviving majority can communicate.
  • Raft election, current-term commit, and current-authority read rules matter beyond counting acknowledgments.
  • Save the job’s accepted ownership version and check it in the same atomic update as the result.
  • An operation ID identifies a repeated action; a fencing token checks whether its worker may still write. You may need both.

Interview tips

  • Draw the actual overlapping sets before applying R + W > N.
  • Pause an old worker, let a new generation publish, then resume the old one and point to the exact rejection.
  • Ask whether a lease or lock protects the actual destination, not just a separate coordinator key.

Important qualifications

  • A maximum-seen fence rejects older writes only after the newer fence reaches that resource; immediate revocation needs a stronger check.
  • Do not equate a controller election term with a per-job ownership epoch.
  • A proven atomic-register protocol can use quorum reads with write-back; consensus is not synonymous with every linearizable register implementation.

Technical references

Practice marks stay in this browser.