When the Replicas Cannot Communicate
In the previous chapter, we replicated Shopend’s database across several servers. One of the replicas, the primary, accepts all application writes and passes its changes to the other replicas. The API instances send their writes to the primary and send suitable reads to the replicas. When the primary fails, failover moves its role to another replica.
In each of those designs, the database servers could communicate with each other. Now suppose they cannot. To keep the example small, assume Shopend has one primary and one replica. Both database servers are working, and every API instance can still reach both of them. But the network connection between the primary and the replica breaks.
Network partitions
Here is the system after the connection breaks:
graph TB
A[Shopend API Instances] -->|Database requests| DB
subgraph DB[Shopend Database Replicas]
direction LR
P[(Primary)] -. Connection broken .- R[(Replica)]
end
The API instances send writes to the primary and reads to the replica, as before. Each database server can still receive requests and run queries against the data it holds. But the primary cannot send its changes to the replica, and the replica cannot send acknowledgments back.
This is called a network partition: parts of the system cannot communicate with each other, even though the servers themselves may still be working. A failed network switch or a misconfigured firewall can cause one. It may last a few seconds or several hours.
A partition is not the same as a server failure. A failed server stops serving requests. During a partition, the servers keep running and keep serving the requests that reach them. Here, the replica keeps answering reads, but it no longer receives the primary’s changes. Each time the primary accepts a write, the replica falls further out of date.
Failover during a partition
Suppose the database promoted a replica whenever that replica could not reach the primary. The replica cannot tell a primary that has crashed from a primary that is still running behind a broken connection. In both cases, no changes arrive, and the primary does not respond.
Failover then promotes the replica. The old primary is still running, and the API instances can still reach it. Nobody has told it that it was replaced, because that message would have to cross the broken connection. So both servers now accept writes as the primary. An API instance that has learned about the new primary sends its writes there. An API instance that has not learned about it keeps sending writes to the old primary.
Consider cart edits, which a primary confirms without waiting for a replica. A shopper changes the quantity of a mug in their cart from 1 to 2, and the request goes to the old primary. Then they remove another item from the cart, and that request goes to the new primary. Both edits succeed. The old primary has the new quantity but still has the removed item. The new primary has removed the item but still has a quantity of 1. Neither server holds the cart the shopper ended up with.
Purchases do not go wrong this way, because we decided that a purchase waits for a replica to acknowledge it. The old primary cannot reach the replica, and the new primary has no replica left, so neither one can confirm a purchase.
The previous chapter said that failover must prevent the old primary from accepting writes, because a server that stops responding has not necessarily stopped running. During a partition, the old primary is still running. To prevent it from accepting writes, something has to reach it, and the replica cannot reach it.
This is why databases use the majority rule from the previous chapter. With one primary and one replica, neither side of the partition has a majority. The replica is not promoted, and the old primary stops accepting writes. Shopend accepts no writes until the connection is restored. Both servers can still answer reads.
This outcome comes from our simplified example, not from Shopend’s design. A real Shopend database would run many more servers. Suppose it runs five servers, and a partition separates three of them from the other two. The three have a majority. They keep their primary, or promote one of themselves if the primary is on the other side, and Shopend keeps accepting writes. The two on the other side cannot have a primary. They can still answer reads, but their data falls behind.