When the Primary Fails

Suppose the primary stops responding while shoppers are placing orders. The API instances are still running, but they cannot complete purchases through it. Sending writes to an arbitrary replica would undo our decision to have one place accept them.

The replicas have copies of the data. We can make one of them take over the primary’s role. Moving that role to another database server is called failover.

Choosing a replacement

The replicas may have made different amounts of progress applying the primary’s changes. Suppose one has applied all changes through a recent purchase, while another is still applying changes from before that purchase. Choosing the second would leave us further behind.

For our design, we select the replica that has applied the most changes among those available and eligible to take over. The database compares how far the replicas have progressed and chooses the replacement. The application does not choose a replica itself.

Being furthest along does not mean a replica has every change the old primary recorded. It means that, among the candidates, it has caught up the most.

Moving the primary’s role

Failover involves three steps:

  1. Select a replica to take over.
  2. Make it the new primary and prevent the old primary from accepting writes.
  3. Direct application writes to the new primary.

Preventing the old primary from accepting writes is part of the change. A server that stops responding has not necessarily stopped running. It might be unreachable because of a network problem, or it might restart later. If it continues accepting writes while the replacement also accepts them, we once again have two copies making independent decisions about purchases.

The other replicas follow the new primary and apply its changes. Database software or its failover tools manage the change of roles, and client libraries can discover the new primary and direct requests to it. Selecting a replacement and enforcing one primary require coordination, but the application does not implement that coordination itself.

What the application experiences

Writes may pause while failover takes place. An API instance may have a request in progress when its connection to the old primary fails, and knowing the new primary’s address does not finish that interrupted request.

The application still has to handle the interruption. In particular, losing a response does not tell it whether a purchase completed. A retry must use the same idempotency key, so that if the purchase is already recorded on the new primary, the server can return that order instead of placing another one.

Replicas let us resume work on another server instead of waiting for the original server to return. They do not make every request succeed during the change, and the replacement can serve only the data it has received.