Why agreeing is difficult
A timeout cannot tell you whether a peer crashed, the network dropped a packet, or a process paused. A write may have committed even though the response never arrived. Build protocols around uncertainty rather than assuming a timeout undoes remote work.
Replication copies state. Consensus chooses a shared order of decisions under a stated fault model. A cluster can replicate asynchronously without providing consensus or linearizable reads.
Majority quorum intuition
In a five-voter cluster, a majority is three. Any two majorities intersect, so conflicting decisions cannot both win if voters follow the protocol. Losing two voters can leave a working majority; losing three cannot. A two-voter cluster needs both and therefore cannot tolerate one unavailable voter for majority progress.
This arithmetic is necessary but insufficient: membership changes, log history and read protocols matter. For quorum-style databases, R + W > N creates overlap under simplified assumptions; it is not a universal proof of linearizability when sloppy quorums, failures or divergent versions are involved.
Experiment with quorum loss
Try one, two and three unavailable voters. Observe when the cluster must stop committing rather than pretending a minority can safely decide. The experiment illustrates arithmetic, while the protocol below provides the additional safety rules.
Raft at interview depth
Raft separates leader election, log replication and safety. Followers receive heartbeats, candidates seek votes after election timeouts, and terms identify election epochs. The leader proposes log entries and replicates them. Rules governing election eligibility and commitment prevent an elected leader from casually discarding committed history.
A partitioned former leader may still believe it leads. It cannot commit new entries without the required quorum. Randomized election timeouts reduce repeated ties. A production implementation also needs snapshots, persistent state and safe membership change.
Consensus generally costs coordination and network round trips; apply it where correct shared decisions matter, not for every disposable analytics counter.
Clocks and causality
Wall clocks can drift or jump. Use a monotonic clock for local elapsed-time deadlines. Lamport clocks preserve an ordering implication for causally related events, but a lower Lamport timestamp does not prove causation. Vector clocks can represent concurrent versions at a metadata cost. Hybrid logical clocks combine physical-time usefulness with logical ordering rules.
For a feed, approximate time order may suffice. For transferring money, timestamp sorting alone cannot establish transactional correctness. Clock uncertainty becomes especially important in lease-based protocols.
Distributed locks are not magic
A local mutex coordinates threads in one process. A distributed lock coordinates participants through an external authority, commonly using leases and renewal. The holder can pause longer than its lease, wake up and continue acting after someone else acquired the lock.
A acquires lease with token 41 -> A pauses
Lease expires -> B acquires token 42 -> B writes
A resumes with token 41 -> storage must reject stale token
A fencing token is a monotonically increasing ownership epoch checked by the protected resource. It works only if that resource enforces the comparison atomically. Merely printing the token in a log does not prevent stale writes. etcd provides revisions, leases and concurrency primitives that support coordination designs; understand their guarantees before building a lock on top.
Decide whether a lock protects cost or correctness
If two workers occasionally refresh the same harmless cached result, a lock may primarily prevent wasted computation. If two workers can both ship an order or overwrite authoritative state, correctness depends on exclusion. Those cases need different failure tolerance.
A worker can pause beyond its lease, resume and send a delayed write after another worker acquires ownership. Renewals and a random ownership token do not make that stale write safe by themselves. The protected store must atomically reject an older fencing epoch after observing a newer one. A lock algorithm and a safe storage-write protocol are separate responsibilities. If the target cannot enforce this, prefer an invariant it can enforce, such as a unique constraint or conditional state transition.
Prefer simpler invariants where possible
For reserving the last inventory unit, an atomic conditional update or database transaction can be safer than a separate distributed lock:
UPDATE inventory
SET available = available - 1
WHERE sku = 'book-17' AND available > 0
RETURNING available;
Zero rows means the reservation failed. Pair the reservation with an idempotent operation identity and appropriate transaction scope. A lock is useful for coordinated work that cannot be expressed in one database operation, but it introduces its own availability and recovery requirements.
Failure detection, split brain and conflict repair
Heartbeats and gossip distribute liveness information; suspicion is not proof. A partition can create different views of membership. Prevent concurrent authoritative writers through quorum/epochs or deliberately permit conflicts with an explicit reconciliation rule.
Last-write-wins can lose valid updates. CRDTs use mathematically designed merge rules for suitable data types; they do not make every business invariant coordination-free. A grow-only set is straightforward; a global stock limit or unique booking is harder.
Interview exercise
Design one scheduled report per tenant per day. Start with a unique database key and idempotent output publication. Explain lease expiry, worker pauses, overlapping runs, fencing and reconciliation. Test a worker that loses its lease immediately before publishing.
Next: transactions and failure recovery.