Replicas
In the replication chapter, we added replicas for two reasons: to spread reads across servers, and to have a server that can take over when the primary fails. Each variant needs replicas for one or both reasons.
t.co and Drive need replicas for both. Their reads outnumber their writes, 100 to 1 for t.co and 10 to 1 for Drive. For t.co, the replicas serve the requests that miss the cache. Drive has no cache, so its replicas serve every follow request.
O’Reilly needs a replica for failover only. Its links must keep working for as long as the books are in print. One replica keeps them working when the primary fails.
O’Reilly also needs backups, because a link must last for years. If someone deletes a link by mistake, the replica deletes it too. A backup still has it.
In every variant, writes still go to the primary.
A new link on a replica
Someone publishes a post with a link. One second later, a different user clicks the short link to reach the original URL. The read goes to a replica that has not applied the new link yet. The replica does not find the code, so the API instance sends 404 Not Found to that user.
We do not want that. The link exists on the primary, and the user who clicked it should get the redirect.
In the replication chapter, we provided read-your-writes by sending a shopper’s reads to the primary for a few seconds after the shopper changed something. That does not help here. The user who wrote the link is the one who posted it. The user who clicks the short link is a different user. So that user’s read still goes to a replica.
So we use a different rule. When a replica does not find a code, the API instance reads the primary. This is like a cache miss with cache-aside: we try the copy first, and on a miss we read the primary. The difference is that we do not write anything back to the replica. Replication adds the link to the replica on its own.
Misses are rare, because almost every code we receive came from a link we created. So these extra reads add little load to the primary.
Guessing codes on Drive
Suppose someone sends many requests with made-up codes to find documents. Each made-up code misses on the replica, and each miss is a read on the primary. So many guessed codes cause many extra reads on the primary.
We rate-limit misses for each client. A client that goes over the limit receives 429 Too Many Requests. We did the same for the Public Courses API.
A network partition
Suppose a replica cannot reach its primary. For each operation and each variant, we choose between consistency and availability. The variants do not need the same choice.
Following a t.co link: AP. The replica keeps redirecting with the links it has, even though its copy may be out of date. A link blocked during the partition can still redirect on that replica. If a new code is missing, the API instance still tries the primary. It cannot reach the primary across the partition.
We accept a blocked link that still redirects for two reasons. Refusing every redirect on an isolated replica would break links across X. Sending all those reads to the primary would exceed its capacity. Choosing availability here is an exception to the blocking deadline: during a partition, a block may take more than a few minutes to reach every replica.
Following a Drive link: CP. Turning off a link must take effect within seconds, so we cannot keep redirecting from an out-of-date replica. So every read of the link’s status must be strongly consistent. Before a replica answers, it must check with the primary that it has the latest committed state. If the replica cannot do that, the API instance reads the primary directly. If neither one can give a consistent read, the request fails. Checking whether the replica’s connection to the primary is alive is not enough: a connected replica can still be behind.
Following an O’Reilly link: AP. We would rather serve an old target for a while than have a printed link stop working. Normally reads go to the primary. During a partition, the replica may take over through failover. Failover stops the old primary from accepting writes. If the replica has not received a link’s target change yet, it serves the old target for that link. Staff must accept that a target change may take longer to reach readers during a partition. This choice is about reads only. It does not let both servers accept writes.
Creating a link: CP. Creating a link is a write, so it waits for the primary. If the primary cannot accept the write, creating the link fails, and the client tries again later. The other choice is AP: the replica accepts the new link. Then two servers accept writes, and their copies disagree until we merge them. This is the problem of the last blue mug, and the reason we keep one primary. For Drive, which uses random codes, the two sides could even give the same code to two different URLs. The cost of CP is small here. The client is the company’s own software, such as X’s servers or Drive’s share dialog, and it can retry. No one who clicks a link is affected.
Turning off a Drive link: CP for the write and later reads. The write goes to the primary. Once it succeeds, every later successful follow request must see the disabled status. If the primary cannot safely accept the write, turning off the link fails and the owner retries later. Changing the document’s sharing setting is a separate operation. It does not change the requirement: a turned-off link must stop working within seconds.