Consensus and Coordination Algorithms Questions
How independent nodes agree on shared state: Paxos and Raft, leader election, quorum reads and writes, distributed locks, coordination services such as ZooKeeper or etcd, and Byzantine fault tolerance for replicas that cannot be trusted. Covers split-brain avoidance, fencing tokens, and the cost of coordination on throughput and latency. Frames when consensus is required versus when it can be designed away.
Design a lease-based leadership mechanism that allows followers to serve local linearizable reads without contacting the leader on every read. Describe how you would handle clock skew, lease renewal, and leader transfer. State assumptions about clock drift or synchronize time service as needed.
Sample Answer
Direct answer
Give selected followers a time-bounded read lease from the leader, and make every write wait until each current lease holder has acknowledged it before the write is reported committed. A follower holding a valid lease then knows it has seen every committed write, so it can answer reads from local state with no leader round trip. Safety rests on three rules: follower leases never outlive the leader's own lease, lease timing uses monotonic clocks with a drift margin, and a new leader cannot commit until every lease granted by the previous one has expired or been revoked.
Terms first
- Linearizable read: the read returns the result of the latest write that completed before the read began, as if there were one copy of the data. A follower that is merely "caught up recently" is not enough; one missed write breaks it.
- Lease: a promise that holds for a fixed time, measured on a clock. Unlike a lock, it expires on its own, so a crashed holder cannot block the system forever.
- Leader lease: in Raft-style systems, a majority of followers promise not to vote for a new leader until a time has passed; while that promise is valid, the current leader knows it is still the only leader.
- Clock drift: two quartz clocks tick at slightly different rates. What matters for leases is the rate error (how fast a duration is measured), not the absolute time of day.
Assumptions
- Each node measures durations with a monotonic clock (one that never jumps backward when NTP, the Network Time Protocol, adjusts the wall clock).
- Clock rate error is bounded by
ε = 0.001(1,000 parts per million). Quartz oscillators are normally far better than this, and the reference NTP daemon limits how fast it slews the clock (nudges it gradually toward the correct time, rather than jumping) to 500 ppm, so 0.001 leaves margin for both. - Process pauses (garbage-collection pauses, VM freezes) are possible and must be handled explicitly, not assumed away.
- Writes go through a Raft log (Raft is a leader-based consensus protocol: the leader appends each write to a replicated log and commits it once a majority has stored it); the leader holds a leader lease renewed via heartbeats (periodic AppendEntries messages, Raft's replication RPC for sending log entries, sent even when there is no new entry to carry, purely to prove the leader is alive).
Design
sequenceDiagram
participant C as Client
participant L as Leader
participant F as Follower (lease holder)
participant G as Other follower
F->>L: request read lease (t0 on F's clock)
L-->>F: grant for T_g, recorded in lease table
C->>L: write x=5
L->>F: AppendEntries(x=5)
L->>G: AppendEntries(x=5)
F-->>L: ack
G-->>L: ack
Note over L: majority acked AND all lease holders acked
L-->>C: write OK
C->>F: read x
Note over F: lease valid, entry applied
F-->>C: 5 (no leader contact)
1. Granting and using a lease
- The leader keeps a lease table: holder, granted duration
T_g, and the time it sent the grant on its own monotonic clock. - Write rule: an entry counts as committed for clients only when (a) a Raft majority has persisted it, and (b) every node with an unexpired lease has acknowledged receiving it. If a lease holder is unreachable, the leader waits until that holder's lease has certainly expired, then drops it from the table and proceeds.
- Read rule on the follower: if the lease is valid now, serve the read from local state at the follower's applied index (the highest log position it has actually executed into its local state; being applied is one step past being committed). If the follower has acknowledged an entry touching that key (told the leader it has persisted it) but has not yet learned it is committed (confirmed by a majority, so it is safe to execute), the read waits for the commit notification, which the leader sends immediately to lease holders as soon as the majority is reached. This closes the gap between ack, commit and apply. Concretely: entry 46 sets
x=9. Follower F persists it and sends its ack at t=10 ms; the leader reaches a majority and marks 46 committed at t=11 ms and tells F immediately; F applies it (updates its in-memoryx) at t=11.2 ms. A read forxthat arrives at F at t=10.5 ms, after the ack but before the commit notice, must wait that roughly 0.5-1 ms rather than read the oldx, or it would return a value one write behind what a client already believes was written.
2. Clock skew: the safety margin
The holder must stop using a lease before the leader stops honouring it, in real time. The holder starts its timer at t0, when it sent the lease request, which is earlier than the grant, so network delay only makes it more conservative. A local duration T takes between T/(1+ε) and T/(1-ε) of real time, depending on whether the clock measuring it runs fast or slow. Each side must assume its own clock is wrong in the direction that hurts it: the holder assumes its clock is running slow, because a slow clock makes a fixed real-time duration look shorter than T on that clock, so waiting for the same clock-reading T_h actually eats into more real time, T_h/(1-\varepsilon); the leader assumes its own clock is running fast, because a fast clock reaches a clock-reading of T_g before as much real time has passed, so its promise really only lasts T_g/(1+\varepsilon) of real time. Using each side's worst case, the holder's lease ends no later than the leader's if
Worked numbers: with T_g = 2000 ms and ε = 0.001,
Drift costs only 4 ms. The real threat is a process pause: a follower checks "lease valid", freezes in a GC pause for 3 seconds, then answers from stale state. The fix is ordering: the follower first takes a local snapshot of the data, then checks that the lease is still valid, and serves the snapshot only if the check passes. The snapshot was therefore current at an instant when the lease was valid, and a pause after the check only delays the reply; it cannot make the answer older than that instant. The holder also uses T_h = 1800 ms rather than 1996 ms, a 196 ms guard for timer imprecision and for a clock that is worse than its assumed bound.
3. Lease renewal
- The follower renews every
T_g / 4 = 500 ms, piggybacked on the heartbeat reply, so a healthy holder never lapses and renewal adds no extra messages. - A renewal restarts both clocks from the new request time; the leader keeps the maximum expiry it has ever granted each holder.
- If renewal fails, the follower falls back to forwarding reads to the leader (or to a ReadIndex check, where the follower asks the leader for its commit index and waits to apply it) until it regains a lease. Correctness degrades to "slower", never to "stale".
4. Leader transfer and failover
- Follower leases never outlive the leader lease. The leader only grants follower leases that expire before its own leader lease does. Because followers refuse to vote for a new leader while the old leader lease is valid, no new leader can exist while any old follower lease is still valid.
- Unplanned failover: the new leader additionally waits out the maximum lease duration (
T_g, stretched by1+ε) from its election before acknowledging any write. Reads during that window are served by the new leader after a ReadIndex round. - Planned transfer: the old leader stops granting and renewing, sends explicit revocations, waits for holders to ack them (fast, one round trip), and only then tells the target to start its election. Planned transfers therefore cost milliseconds, not the full lease.
Worked example of the cost
Five nodes, three followers hold read leases, T_g = 2 s. Normal case: a write waits for the slowest of the three lease holders instead of the fastest two of four followers, so write latency tracks the slowest lease holder. If one lease holder's machine dies, writes stall until its lease is certainly expired: up to 2000 × 1.001 ≈ 2002 ms after its last renewal. That is why lease holders should be few (the replicas that actually serve read-heavy traffic, for example one per region) and why T_g is seconds, not minutes.
Trade-offs and pitfalls
| Choice | Gain | Cost |
|---|---|---|
| Follower read leases (this design) | Local linearizable reads, zero leader round trips | Writes wait for all lease holders; a dead holder stalls writes up to T_g |
| ReadIndex on follower | No clock assumptions | One leader round trip per read (batchable) |
| Bounded-staleness follower reads | Cheapest | Not linearizable |
Pitfalls a senior answer names:
- Using the wall clock (
gettimeofday) instead of a monotonic clock: an NTP step backward silently extends a lease. - Starting the holder's timer when the grant arrives instead of when the request was sent: network delay then eats the safety margin.
- Letting follower leases outlive the leader lease: a new leader commits a write that an old lease holder never sees.
- Ignoring VM live migration (a virtual machine moved between physical hosts while still running, which freezes it for a window similar to a GC pause) and long GC pauses: they are the realistic way a "valid" lease is used after expiry.
Provide a high-level client library interface (pseudo-code) for a distributed lock client with operations acquire(key, ttl), release(key), and renew(key, ttl) backed by a consensus store (e.g., etcd). Discuss failure modes (client crash, network partition), lease expiry, clock skew, and safe usage patterns you'd document for teams.
Sample Answer
Short answer
The client wraps three primitives of etcd (a strongly consistent key-value store replicated with the Raft consensus algorithm). The ttl argument is a time-to-live: how long the lock survives without renewal. A lease is a server-side timer the client must keep refreshing; if it stops, the lease expires and every key attached to it is deleted. acquire(key, ttl) grants a lease and creates the lock key only if it does not exist, in one atomic transaction, and returns a handle containing the lease ID and a fencing token (a number that grows with every successful acquire, here the key's create revision: etcd keeps one cluster-wide counter that increases by one on every write anywhere in the store, and remembers that counter's value at the moment each key was created, so this number only ever goes up). renew refreshes the lease. release deletes the key only if it is still the one this handle created. The part that matters most for teams is the usage contract: a lease-based lock can expire while its holder is paused or partitioned, so callers must check renewal results, stop work when renewal fails, and pass the fencing token to the resource they protect.
Interface (pseudo-code)
struct LockHandle { key, leaseId, fence, deadline } // deadline measured on a monotonic clock
acquire(key, ttl, waitUpTo) -> LockHandle or TIMEOUT
start = monotonicNow()
leaseId = store.leaseGrant(ttl)
loop:
txn if createRevision(key) == 0 // key absent
then put(key, ownerId, lease = leaseId)
if txn.succeeded:
fence = txn.revision // create revision of our key
return LockHandle(key, leaseId, fence, start + ttl - safetyMargin)
if monotonicNow() - start > waitUpTo:
store.leaseRevoke(leaseId); return TIMEOUT
watch(key) until DELETE or remaining wait elapses // "watch": block until the server pushes a notification that this key changed, instead of polling
renew(handle, ttl) -> bool
sent = monotonicNow()
remaining = store.leaseKeepAlive(handle.leaseId) // server returns remaining TTL
if remaining <= 0: return false // lease already gone: we lost the lock
handle.deadline = sent + remaining - safetyMargin // measured from BEFORE the request
return true
release(handle) -> bool
txn if createRevision(handle.key) == handle.fence // still OUR key, not a successor's
then delete(handle.key)
store.leaseRevoke(handle.leaseId)
return txn.succeeded
isHeld(handle) -> bool
return monotonicNow() < handle.deadline
The library runs renewal on a background timer at roughly TTL/3 and exposes isHeld() plus a lost-lock callback. etcd itself ships a similar recipe in its concurrency package (a session holding one lease, plus a mutex built on it) that production Go code should prefer over a hand-rolled version.
Executed demonstration against a real etcd
This drives etcd's JSON gateway (etcd's built-in HTTP+JSON API, a thin translation layer over its native gRPC interface, so plain HTTP calls can drive it with no generated client) directly (standard library only, no client package), so it can be run as is against a throwaway single-node etcd. Start one with docker run -d --name etcd-demo -p 2379:2379 quay.io/coreos/etcd:v3.5.13 etcd --listen-client-urls http://0.0.0.0:2379 --advertise-client-urls http://127.0.0.1:2379 and remove it afterwards with docker rm -f etcd-demo.
import base64, json, os, time, urllib.request, uuid
ETCD = os.environ.get("ETCD", "http://127.0.0.1:2379")
b64 = lambda s: base64.b64encode(s.encode()).decode()
def call(path, body):
req = urllib.request.Request(ETCD + path, json.dumps(body).encode(),
{"Content-Type": "application/json"})
return json.load(urllib.request.urlopen(req))
def acquire(key, ttl_s):
lease = call("/v3/lease/grant", {"TTL": ttl_s})["ID"]
owner = uuid.uuid4().hex
# /v3/kv/txn runs one atomic if/then/else on the server: IF every entry in
# "compare" holds, run "success"; otherwise run "failure". Nothing else can
# run between the check and the write, which is what makes this safe.
resp = call("/v3/kv/txn", {
# succeed only if the key does not exist (create_revision == 0)
"compare": [{"key": b64(key), "target": "CREATE", "create_revision": "0"}],
"success": [{"request_put": {"key": b64(key), "value": b64(owner), "lease": lease}}],
"failure": []})
if not resp.get("succeeded"):
call("/v3/lease/revoke", {"ID": lease})
return None
fence = int(resp["header"]["revision"]) # revision at which we created the key
return {"lease": lease, "owner": owner, "fence": fence}
def renew(h):
ttl = call("/v3/lease/keepalive", {"ID": h["lease"]}).get("result", {}).get("TTL")
return ttl is not None and int(ttl) > 0
def release(key, h):
resp = call("/v3/kv/txn", {
"compare": [{"key": b64(key), "target": "CREATE", "create_revision": str(h["fence"])}],
"success": [{"request_delete_range": {"key": b64(key)}}], "failure": []})
return bool(resp.get("succeeded"))
key = "/locks/billing-reconcile"
a = acquire(key, 2)
print("A acquired:", a is not None)
print("B acquire while A holds:", acquire(key, 2) is not None)
print("A renew while alive:", renew(a))
time.sleep(4) # A stops sending keepalives (crash or partition)
print("A renew after lease expired:", renew(a))
b = acquire(key, 10)
print("B acquired after expiry:", b is not None)
print("B fence > A fence:", b["fence"] > a["fence"])
print("A release (stale handle):", release(key, a))
print("B release:", release(key, b))
Output:
A acquired: True
B acquire while A holds: False
A renew while alive: True
A renew after lease expired: False
B acquired after expiry: True
B fence > A fence: True
A release (stale handle): False
B release: True
What this shows: mutual exclusion while A's lease is alive; the lock freeing itself when A stops renewing (the crash case); A's stale release failing because the compare is on its own create revision, so it cannot delete B's lock; and B's fencing token being strictly greater than A's.
Failure modes
| Failure | What happens | What the library does |
|---|---|---|
| Client crash | Keepalives stop, lease expires after at most TTL, key deleted, next waiter's watch fires | Nothing to do: this is the design working |
| Client partitioned from etcd | Same as a crash from etcd's view. The client may still be running | Renewals fail; library fires lost-lock callback when deadline passes, even if the network is silent |
| Client paused (garbage collection, VM freeze) longer than TTL | Lease expires, another client acquires, paused client wakes up believing it holds the lock | Unfixable in the client. Only a fencing check at the resource rejects its late writes |
| etcd leader change | Brief write unavailability; etcd extends lease deadlines on leader change so sessions do not all expire at once | Retry renewals with backoff (waiting a bit longer between each retry instead of hammering immediately again) within the deadline |
| etcd minority partition | Members without quorum (a majority of the cluster's voting members) cannot grant, renew, or commit | Client treats renewal failure as lock loss |
Lease expiry and clock skew
The lease TTL is enforced by the etcd leader's timers, not by client wall clocks, so two clients with skewed clocks never disagree about who holds the key. Skew and timing still bite in two places:
- The client's own estimate. The client must compute "I still hold it until" from a monotonic clock (one that never jumps when NTP, the Network Time Protocol, corrects the system time), started before sending the renewal, minus a safety margin. Measuring from when the reply arrived overestimates the remaining time by the network delay.
- Rate differences. If the client's clock runs slower than the leader's, the client believes it has more time than it does. The safety margin (for example 10 to 20% of TTL) absorbs normal drift, but no margin covers a multi-second pause, which is again why fencing exists.
Safe-usage guidance I would document for teams
- Always use
withLock(key, ttl, fn), which acquires, starts renewal, cancelsfn's context (Go's standard mechanism for telling a running function "stop, you are no longer allowed to continue") on lock loss, and releases infinally. - Pass
handle.fenceto every write to the protected resource, and have the resource reject tokens lower than the highest it has seen (for a SQL table:UPDATE ... WHERE last_fence <= :fence). - Never do unbounded work between
isHeld()checks; check before each side-effecting step. - Pick TTL at least several times the renewal interval and longer than your worst observed pause, but short enough that recovery after a crash is acceptable (10 to 30 s is a common range; state yours).
- Acquire multiple locks in sorted key order, and always with a wait deadline.
- Do not use the lock for anything whose side effects cannot be fenced or made idempotent (safe to run more than once, because repeating it produces the same result rather than a duplicate side effect) without also accepting that they may occasionally run twice.
Contrast: DynamoDB conditional writes (non-consensus from the caller's view)
Amazon DynamoDB (a managed key-value database) can implement the same pattern without running a consensus cluster yourself: acquire with PutItem and a condition expression such as attribute_not_exists(lockKey), renew and release with conditional updates that check an owner or version attribute. Differences a team should know:
- No leases or watches. Expiry must be implemented by the client. DynamoDB's TTL feature is a clean-up job that deletes expired items "typically within a few days", so it cannot be the lock timeout.
- Clock handling. Storing an absolute expiry time makes safety depend on client clocks agreeing. AWS's open-source DynamoDB Lock Client avoids that: it stores only a relative lease duration plus a record version number (a GUID, a long random identifier, regenerated on every heartbeat so a waiter can tell the record actually changed rather than staying identical by coincidence), and a waiter declares the lock stale only if the version has not changed for a full lease duration measured on its own clock.
- Waiting is polling. Without watches, waiters poll, which costs read capacity (DynamoDB's read-throughput budget, consumed by every poll whether or not anything actually changed) and adds latency.
- Fencing can use a counter attribute incremented atomically on each acquire.
Choose DynamoDB when you are already on AWS, lock traffic is modest, and you do not want to operate etcd. Choose etcd when you need push notifications (watches), fast failover detection via leases, or you already run it (for example in Kubernetes, a container-orchestration platform that already runs its own etcd cluster internally).
Complexity and edge cases
Each operation is one round trip to the store plus a quorum write (acquire, release) or a leader-handled keepalive (renew); waiting is event-driven via a watch, not polling. Edge cases to handle: acquire succeeding on the server but the reply being lost (the lease then expires on its own because the client never renews it; or reuse a client-chosen owner ID so a retry can recognise its own key), release after the lease already expired (returns false, not an error), and double release (idempotent).
What is split-brain in the context of a replicated cluster? Describe typical causes, the consequences for correctness and availability, and how you would prevent or mitigate it in a production system.
Sample Answer
Short answer
Split-brain is when a replicated cluster that is supposed to have one leader (or one writable primary) ends up with two or more nodes that each believe they are in charge, usually because a network failure split the cluster into groups that cannot see each other. If both sides keep accepting writes, the data forks: the same record can be changed two different ways, and there is no longer one correct version. The standard prevention is to require a majority quorum (more than half the voting members) before anyone may act as leader, so at most one side of any split can qualify, and to fence a deposed leader so its late writes are rejected.
Typical causes
| Cause | What happens |
|---|---|
| Network partition plus failover without quorum | Each side sees the other as dead; a standby promotes itself while the old primary keeps serving |
| Even-sized cluster split in half | With 4 nodes split 2/2, a "majority of what I can see" rule lets both halves elect a leader |
| Process pause (garbage collection: a runtime's automatic memory cleanup, which can briefly stop the whole process; or a VM freeze) | The leader stops for longer than the failure timeout, a new leader is elected, and the old one wakes up and writes as if nothing happened (a "zombie" leader) |
| Clock-based leases (a lease is a time-boxed grant of leadership that expires automatically unless renewed) with clock skew | The old leader thinks its lease is still valid while followers think it has expired |
| Asynchronous replication (the primary confirms a write to the client before every replica has stored it) with automatic promotion | Primary and replica both writable after a hiccup; also loses writes that had not yet reached the replica when it was promoted |
| Human error | Someone manually promotes a replica during an incident without demoting the primary |
Consequences
For correctness:
- Conflicting updates to the same key (two different balances, two owners of one seat).
- Lost updates when one side's writes are later discarded during recovery.
- Violated invariants that a single leader would have enforced: uniqueness, "stock never below zero", "each job runs once".
- Duplicate side effects: both leaders send the same email, charge, or event.
For availability, the design choice is stark. During a partition a system can either let both sides serve (available but inconsistent, the split-brain outcome) or let only the majority side serve (consistent, but the minority side refuses writes). This is the CAP theorem trade-off (consistency, availability, partition tolerance: under a partition you must give up one of the first two) in concrete form. Quorum-based systems deliberately choose to make the minority unavailable.
Worked example: why node count matters
Majority means strictly more than half.
- 5 nodes, partition into 3 and 2. Majority is 3. The 3-node side elects a leader and keeps working; the 2-node side cannot, so it refuses writes. Exactly one leader.
- 4 nodes, partition into 2 and 2. Majority is 3. Neither side has it, so with a correct quorum rule neither side can write: no split-brain, but a full write outage. With a broken rule ("elect if you can reach half"), both sides elect: split-brain.
- 2 nodes (primary plus replica). Any partition is 1 and 1. There is no majority possible, so automatic failover between two nodes is inherently unsafe without a third party. This is why two-node high-availability setups add a witness (a small voting member, often in a third location, that stores no data but breaks ties).
This is why clusters are sized 3 or 5: they tolerate 1 or 2 failures respectively, and an even-sized cluster tolerates no more failures than the odd size just below it.
Prevention and mitigation in production
- Majority quorum for leadership. Use a consensus protocol (Raft, Paxos, or Zab, the protocol inside the ZooKeeper coordination service) or a store built on one to decide who leads. A leader that loses contact with a majority must step down (Raft implementations call this CheckQuorum).
- Odd number of voters across failure domains (independent zones, such as separate data centers or power feeds, that are unlikely to fail at the same time). For example 3 nodes in 3 zones, or 2 + 2 + 1 across three sites, so no single site failure leaves the survivors without a majority.
- Fencing. Each new leader gets a higher epoch or term number (a counter that only ever goes up; whichever side shows the higher number is the more recent leader), and storage rejects writes with an older one. This handles the zombie leader that wakes up after a pause. Where storage cannot check epochs, STONITH ("shoot the other node in the head") forcibly powers off or isolates the old primary before the new one starts.
- Lease discipline. If leaders use time-based leases, the leader must treat its lease as expired earlier than followers will, with a margin that covers clock drift.
- Detection. Alert when more than one node reports itself as writable or leader. Split-brain that is detected in one minute is an incident; detected in a day it is a data-recovery project.
- Testing. Inject partitions and pauses in staging (for example with a chaos tool: software that deliberately triggers failures, such as killing a process or cutting network access, to test how the system responds) and verify the old leader actually stops writing.
Trade-offs and pitfalls
- A longer failure-detection timeout makes false failovers rarer but does not prevent split-brain; only quorum plus fencing does. It also lengthens every real outage.
- Quorum means the minority side goes read-only or down during a partition. That is the price of correctness, and it should be a documented, accepted behaviour.
- Split-brain is not only a database problem: schedulers, lock services and leader-elected workers all have the same failure if their leadership is not quorum-backed and fenced.
Implement the Bully algorithm leader election in Python. Given a list of node IDs and a way to 'send' an election message to higher-ID nodes (synchronous, reliable for this exercise), write a function elect_leader(node_id, all_node_ids) -> leader_id that follows the Bully algorithm assumptions and returns the elected leader.
Sample Answer
Direct answer
The Bully algorithm elects the highest-ID live node as leader. A node that starts an election sends ELECTION to every node with a higher ID. If none of them answers, it wins and announces itself with COORDINATOR. If any higher node answers OK, the starter backs off, because that higher node now runs its own election, and the process repeats until the highest live ID finds nobody above it. Under the exercise's assumption (synchronous, reliable delivery, so silence means "crashed"), the function below always returns the highest live ID.
Approach
- Leader election means getting a group of nodes to agree on exactly one coordinator. Bully's rule is simple: biggest live ID wins, and a bigger node "bullies" smaller ones out of the race.
- The question's signature has two parameters plus "a way to send"; here that way to send is a third parameter. If the interviewer wants exactly
elect_leader(node_id, all_node_ids), bind the network withfunctools.partialor a closure. - The "send" is injected as a callable,
send_election(sender, receiver) -> bool, returningTruewhen the receiver replies OK. This keeps the algorithm testable without a network. - Rather than spawning a parallel election for every responder, the code follows the lowest responder. With reliable delivery every responder's election converges on the same winner, so this deterministic walk returns exactly what the full algorithm would, while a second helper counts the messages of the textbook "every responder starts its own election" version.
Code
def elect_leader(node_id, all_node_ids, send_election):
"""Bully election started by node_id.
send_election(sender, receiver) returns True if receiver replies OK
(it is alive) and False if it never replies (it has crashed). The
exercise makes delivery synchronous and reliable, so silence = crash.
Returns the elected leader: the highest ID that is alive.
"""
ids = set(all_node_ids)
if node_id not in ids:
raise ValueError(f"unknown node {node_id}")
current = node_id
while True:
higher = sorted(n for n in ids if n > current)
responders = [n for n in higher if send_election(current, n)]
if not responders:
# No higher node answered, so current wins. In a real system it
# now broadcasts COORDINATOR(current) to every lower ID.
return current
# Someone higher is alive: current backs off. Each responder runs its
# own election; with reliable delivery they converge on one winner,
# so following the lowest responder reaches the same answer.
current = min(responders)
def cascade_message_count(starter, all_node_ids, alive):
"""Messages in the textbook version, where EVERY responder starts its own
election (ELECTION + OK + final COORDINATOR broadcast)."""
started, frontier, election, ok = set(), [starter], 0, 0
while frontier:
n = frontier.pop()
if n in started:
continue
started.add(n)
for h in (x for x in all_node_ids if x > n):
election += 1
if h in alive:
ok += 1
frontier.append(h)
winner = max(started)
coordinator = sum(1 for x in all_node_ids if x < winner)
return winner, election, ok, coordinator
if __name__ == "__main__":
nodes = [1, 2, 3, 4, 5, 6]
def network(alive):
return lambda sender, receiver: receiver in alive
scenarios = [
("all alive, node 2 starts", {1, 2, 3, 4, 5, 6}, 2),
("6 crashed, node 2 starts", {1, 2, 3, 4, 5}, 2),
("5 and 6 crashed, node 1 starts", {1, 2, 3, 4}, 1),
("only the starter alive", {3}, 3),
("starter is highest live ID", {1, 2, 4}, 4),
]
for label, alive, start in scenarios:
print(f"{label}: leader={elect_leader(start, nodes, network(alive))}")
for n in (6, 10, 50):
ids = list(range(1, n + 1))
w, e, o, c = cascade_message_count(1, ids, set(ids))
print(f"n={n}, lowest ID starts, all alive: winner={w} "
f"ELECTION={e} OK={o} COORDINATOR={c}")
Output:
all alive, node 2 starts: leader=6
6 crashed, node 2 starts: leader=5
5 and 6 crashed, node 1 starts: leader=4
only the starter alive: leader=3
starter is highest live ID: leader=4
n=6, lowest ID starts, all alive: winner=6 ELECTION=15 OK=15 COORDINATOR=5
n=10, lowest ID starts, all alive: winner=10 ELECTION=45 OK=45 COORDINATOR=9
n=50, lowest ID starts, all alive: winner=50 ELECTION=1225 OK=1225 COORDINATOR=49
Trace of "5 and 6 crashed, node 1 starts": node 1 sends ELECTION to 2, 3, 4, 5, 6; nodes 2, 3, 4 answer OK. Following node 2: it contacts 3, 4, 5, 6 and hears from 3 and 4. Node 3 hears from 4. Node 4 contacts 5 and 6, hears nothing, and wins.
Key points
- Assumptions the algorithm relies on: every node knows every other node's ID; IDs are unique and totally ordered; delivery is synchronous, so a node that does not answer within a known bound really has crashed. Those are exactly the assumptions the exercise grants.
- The winner broadcasts COORDINATOR to all lower IDs so they learn the result; the code returns the winner instead of simulating the broadcast.
- A recovering node with a higher ID starts an election as soon as it comes back and takes over, even if the current leader is healthy. That is where the name comes from.
Complexity
- Messages, worst case: the lowest ID starts and everyone is alive. Each node
ksends ELECTION to then - knodes above it, for a total of
ELECTION messages, the same number of OK replies, and n - 1 COORDINATOR messages. The printed counts match: 15 for n = 6, 45 for n = 10, 1225 for n = 50. So the textbook algorithm is O(n²) messages.
- Messages, best case: the highest live node starts and sends
n - 1COORDINATOR messages (plus ELECTION probes to any dead higher IDs). - The function itself: each loop iteration moves
currentstrictly upward, so it runs at mostniterations of up tonsends,O(n²)time,O(n)extra space.
Edge cases
| Case | Behaviour |
|---|---|
| Starter is the highest live ID | No higher node answers; it returns itself immediately. |
| Only the starter is alive | Same: returns the starter. |
| Starter ID not in the list | Raises ValueError rather than electing silently. |
| Duplicate IDs in the input | Collapsed by set; in a real system duplicate IDs break the algorithm, so they must be rejected at configuration time. |
| Higher node crashes after sending OK | Not covered by the exercise's reliable model. A real implementation sets a timer after receiving OK; if no COORDINATOR arrives, the node restarts the election. |
Why real systems rarely use Bully
- In an asynchronous network a slow node looks exactly like a dead one, so two nodes can each believe they are leader (split brain, two leaders acting at once). Bully has no quorum and no notion of a term or epoch (a counter that marks which leadership is newer), so it cannot stop the old leader from acting.
- It always prefers the highest ID, so a flaky high-ID node causes repeated leadership changes.
- Production systems use a quorum-based protocol (Raft and ZooKeeper, consensus systems where a leader is elected and writes are accepted only once a majority of servers agree; etcd, a small Raft-based key-value store, offers the same kind of lease) that elects a leader only with a majority and tags every leadership with a monotonically increasing term that others can use to reject a stale leader.
Across three active regions, build a booking coordination protocol that guarantees no double-booking of the same resource while keeping user-facing latency at p95 under 300ms. Weigh consensus (Paxos or Raft), partitioned ownership, optimistic concurrency with conflict resolution, and hybrid approaches against each other on latency, complexity, and availability.
Sample Answer
Direct answer
I would use partitioned ownership: every bookable resource (a room, a seat, a table slot) has exactly one home region, and a booking for it can only commit there, as a conditional write (a write the store only accepts if a stated condition holds, such as "no row already exists for this slot", so two conflicting writes can never both succeed) in a store that replicates with a majority quorum across that region's availability zones. Requests from users in other regions are forwarded, costing one cross-region round trip. Global consensus is kept for the small ownership map (which region owns which resources), not for bookings. Optimistic multi-region writes are ruled out, because they find a double booking only after both customers have been told "confirmed". A single global Raft group does meet 300 ms when healthy, but my numbers below show it breaks the budget (305 ms) after losing its leader region, and it puts every booking on one leader. For fungible inventory ("any of 40 standard rooms") I would add escrow, splitting the count into per-region quotas so most bookings commit locally.
Terms used below
- p95 latency: the 95th percentile. 95% of requests finish faster than this value.
- Raft / Paxos: consensus protocols. A group of replicas agrees on one ordered log of writes, and a write is committed once a quorum (a majority, such as 2 of 3) has stored it. Two majorities always share a member, so two conflicting writes cannot both commit.
- Availability zone (AZ): an isolated datacenter within one cloud region, a few milliseconds away from its siblings.
- Optimistic concurrency control (OCC): write without locking, check for conflicts afterwards, and resolve them (retry, merge or cancel).
- Linearizable: all readers and writers see one order of operations consistent with real time. "Is slot 7 free?" followed by "book slot 7" can then be made atomic.
- Idempotency key: a client-supplied ID for a request, so a retry after a timeout does not create a second booking.
- Fencing epoch: a number that increases each time ownership of a resource moves, checked on every write so a region that has lost ownership cannot keep writing.
Assumptions (stated, so the arithmetic can be checked)
- Regions: us-east, us-west, eu-west. Planning RTTs (round-trip times, round figures rather than measurements): east to west 65 ms, east to eu 75 ms, west to eu 140 ms.
- 25 ms of application and storage work per booking; a quorum commit across 3 AZs inside one region takes about 4 ms.
- Figures below are the server-side part of a single booking. The user's last mile (the final hop from their device to the network, usually the slowest and most variable part of the trip) to their nearest region, plus queueing jitter, has to fit into what is left of the 300 ms, so the design with more headroom is the safer p95 choice.
The four options compared
| Option | Double booking prevented? | Server-side latency (best / worst) | Availability | Complexity |
|---|---|---|---|---|
| Consensus: one global Raft group (leader us-east, one voter per region) | yes: every booking is one linearizable conditional write | 90 / 165 ms healthy; 165 / 305 ms after us-east is lost | survives loss of one region, but every booking depends on one leader and slows sharply after a failover | low for the application, high for operations (leader elections over the wide-area network, or WAN, and a single write bottleneck) |
| Partitioned ownership, 3 AZs in the home region | yes: one owner per resource, conditional write there | 29 / 169 ms | a region outage pauses bookings for resources homed there; other resources are untouched | moderate: routing plus an ownership map |
| Optimistic concurrency with conflict resolution (write in any region, replicate asynchronously) | no: two regions can confirm the same slot, and the conflict is found after replication | about 29 ms everywhere | highest | high: conflict detection, compensation (refunding or rebooking the customer who loses the conflict), customer apology flows |
| Hybrid: partitioned ownership plus escrow for fungible stock | yes: a region can only sell units in its quota | 29 ms for escrowed stock; specific items as in partitioned ownership | quotas keep selling locally during a partition, up to what each region holds | moderate to high: quota rebalancing |
The latency columns come from this calculation:
# Planning assumptions (round figures, not measurements)
RTT = {frozenset({"east", "west"}): 65, frozenset({"east", "eu"}): 75, frozenset({"west", "eu"}): 140}
def rtt(a, b):
return 0 if a == b else RTT[frozenset({a, b})]
REGIONS = ["east", "west", "eu"]
WORK = 25 # app + storage engine work per booking, ms
AZ_COMMIT = 4 # quorum commit across 3 availability zones in one region, ms
def nearest_other(r, alive=REGIONS):
return min(rtt(r, o) for o in alive if o != r)
def global_raft(user, leader, alive=REGIONS):
# forward to the leader, leader waits for one more voter (majority of 3)
return rtt(user, leader) + nearest_other(leader, alive) + WORK
def home_shard_in_region(user, home):
return rtt(user, home) + AZ_COMMIT + WORK
def home_shard_2_2_1(user, home):
# 5 voters: 2 in home, 2 in home's nearest region, 1 elsewhere; majority 3
# needs the leader's partner AZ plus one voter in the nearest other region
return rtt(user, home) + nearest_other(home) + WORK
rows = [
("global Raft, leader east", lambda u, h: global_raft(u, "east")),
("global Raft, east lost, leader west", lambda u, h: None if u == "east" else global_raft(u, "west", ["west", "eu"])),
("home shard, 3 AZs in home region", home_shard_in_region),
("home shard, 2+2+1 across regions", home_shard_2_2_1),
]
for name, f in rows:
vals = [f(u, h) for u in REGIONS for h in REGIONS]
vals = [v for v in vals if v is not None]
print(f"{name:<38} best {min(vals):>4} ms worst {max(vals):>4} ms")
Output (executed):
global Raft, leader east best 90 ms worst 165 ms
global Raft, east lost, leader west best 165 ms worst 305 ms
home shard, 3 AZs in home region best 29 ms worst 169 ms
home shard, 2+2+1 across regions best 90 ms worst 240 ms
Two things to notice. First, in the healthy case the worst paths of global Raft (165 ms) and in-region ownership (169 ms) are almost equal, so latency alone does not decide between them. What decides is the failure case: with us-east gone, a user in eu-west booking through a leader in us-west pays 140 ms to reach the leader plus 140 ms for the only remaining voter to acknowledge, plus 25 ms of work, which is 305 ms before any last-mile time. Second, partitioned ownership commits a resource's bookings in ~29 ms whenever the user is in the resource's home region, which for hotels and restaurants (booked mostly by people nearby or in the same market) is a large share of traffic.
Recommended design
flowchart LR
U[User in any region] --> R[Nearest region API]
R -->|lookup owner| M[Ownership map<br/>global Raft, rarely written]
R -->|home is here| S[Home shard<br/>3-AZ quorum store]
R -->|home is elsewhere| F[Forward to home region API]
F --> S2[Remote home shard<br/>3-AZ quorum store]
-
Ownership map. Resource to home region, plus the current fencing epoch. It changes only when a region is added or ownership is moved, so a global consensus group (for example etcd, which is itself a Raft-replicated key-value store, or another consensus-backed database) costs nothing on the booking path; every region caches it.
-
Booking write, in the home shard. One statement that both checks and books, for example a table with a uniqueness constraint on
(resource_id, slot):sqlINSERT INTO bookings (resource_id, slot, booking_id, owner_epoch, idempotency_key) VALUES ($1, $2, $3, $4, $5); -- unique (resource_id, slot) rejects a second booking of the same slot; -- unique (idempotency_key) turns a client retry into a no-op; -- the shard rejects writes whose owner_epoch is below its current epoch.The store replicates that write to a majority of three AZs before acknowledging, so a single-AZ failure loses nothing and cannot reopen the slot.
-
Holds for checkout flows. "Reserve for 10 minutes while the user pays" is the same insert with
status = 'held'and an expiry; a background sweeper (a process that periodically scans for and cleans up expired records) releases expired holds, and confirmation is a conditional update fromheldtoconfirmed. -
Escrow for fungible stock. If 40 identical rooms are on sale, give us-east 20, us-west 10 and eu-west 10 as local quotas. Like any other resource, the room type itself has one home region recorded in the ownership map; that home region owns the authoritative total (40) and is the only one allowed to change how it is split, while every region, including the home region, decrements its own local quota with a local conditional write. The sum of the local quotas can never exceed 40. A region that runs low asks the home region to move quota from another region's allotment to its own, which is the only cross-region step and is off the booking's critical path.
-
Region loss. Bookings for resources homed in the lost region pause (fail closed) rather than fail over automatically. Moving ownership requires bumping the epoch in the ownership map, and it is only safe once every confirmed booking from the old home is known to the new one. With asynchronous cross-region replication, an automatic failover can lose the last few confirmed bookings and then sell those slots again: a double booking caused by the failover itself.
Why not the others
- Global Raft: correct and simple for the application, but every booking in the world queues behind one leader, a leader election over the WAN can take seconds, and the post-failover worst case (305 ms server-side) already exceeds the budget.
- Optimistic concurrency with conflict resolution: it cannot meet the stated guarantee. If eu-west and us-west each accept slot 7 at the same moment, both users see "confirmed", and the resolver can only cancel one afterwards. That fits businesses that deliberately overbook and compensate (airlines), not a "no double booking" requirement.
- Consensus per booking across regions is really the global-Raft row: paying a WAN quorum for every write gives up the locality that most bookings have.
What would change the recommendation
- "Bookings must continue during a full region outage": place each home shard as 5 voters, 2 in the home region, 2 in its nearest region and 1 in the third. A majority then always includes a voter outside the home region, so a region loss neither loses confirmed bookings nor stops writes. The cost is the 240 ms server-side worst case above, which leaves about 60 ms for the last mile and jitter. I would only choose it after measuring real p95 RTTs.
- Mostly non-local traffic (for example, international travellers booking another continent): locality buys less, and the choice between global Raft and ownership rests on failure isolation rather than latency.
- No specific-item inventory at all (every unit interchangeable): lean fully on escrow, and bookings become local writes almost everywhere.
Trade-offs and pitfalls
- Check-then-write in application code ("SELECT free slots, then INSERT") without a constraint or a conditional write double-books under concurrency, even inside one region. The uniqueness check must be enforced by the store.
- Asynchronous replication presented as high availability. A failover to an asynchronous replica can lose confirmed bookings, which reopens the slots. Say explicitly which replicas are synchronous.
- Hot resources. A concert on-sale sends every request to one home shard. Mitigate with escrow for interchangeable seats and a queue in front of the sale, not by adding writers.
- Retries without idempotency keys. A timeout after a successful commit, retried, becomes a second booking of a different slot for the same customer. The key makes the retry return the original result.
- Budgeting with averages. The p95 includes jitter and queueing. Leave headroom (the recommended design leaves over 130 ms server-side) and verify with measured latency, not planning RTTs.
Unlock Full Question Bank
Get access to all 36 Consensus and Coordination Algorithms interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.