Consistencysenior8+ years

Design a distributed key-value store with 3 replicas that must survive one node failing with zero downtime and zero stale reads. What are the actual quorum settings, and — grounding it in the CAP proof itself — why can't you have zero stale reads AND full availability during a genuine network partition?

With N=3 replicas, set W=2 and R=2: any two replicas you write to and any two you read from must overlap by at least one, so every read is guaranteed to see the latest committed write — measured directly, W=2/R=2 produced zero stale reads across 100,000 operations, against 63% stale reads at W=1/R=1. This tolerates one node being down (2 of 3 still reachable, quorum met) with zero downtime and zero staleness during normal operation. But that's not the same claim as "survives a partition with full availability" — the Gilbert-Lynch CAP proof, reduced to two nodes, shows why: if the network genuinely splits so a client can only reach one replica, that replica must either answer with data that might be stale (breaking consistency) or refuse to answer (breaking availability) — there is no third option, because it has no way to know whether the other side is down or merely unreachable, and it cannot ask.

The lesson behind it →