Sharding the Links
In our design, the replicas and the t.co cache spread the reads. Every write still goes to the one primary. Each replica holds a full copy of the links, so adding replicas does not reduce how much data each server stores. To spread the writes and the storage across servers, we shard the links.
t.co needs sharding. It creates about 2,000 links per second and stores about 36 TB a year.
Drive needs sharding. Its writes fit on one primary, but it stores about 9 TB a year, and links stay for years.
O’Reilly does not need sharding. It creates under one link per second and stores about 36 MB a year, so one server holds its links for as long as the books are in print.
The shard key
Following a link is the most frequent operation, and it looks up one code. So we shard by code. Each follow then goes to one shard.
For Shopend, we sharded by shop. That keeps a shop’s products, carts, and orders on one shard, so a purchase stays within one shard. A link has no group like that to keep together.
Operations across shards
With the links sharded by code, two operations must ask every shard:
-
Drive: turning a document’s links off or back on. The lookup is by document, and that document’s links can be on any shard.
-
t.co: blocking the links to a harmful site. The lookup is by site.
Both are rare compared with following a link. So we ask every shard and combine the results. This is the same pattern as Shopend’s report that counted orders across all shops.
If turning links off and on became frequent, we could add a second table that maps a document to its codes, sharded by document. Turning off a document’s link would then read that table on one shard to find the codes, and update each link on the shard that holds it. For a document with one link, that is two shards instead of every shard. But then every change would have to update both tables, so they stay consistent.
Computing the shard
For Shopend, the router uses a directory that records which shard holds each shop. Shopend has 10,000 shops, so the directory is small. Each router keeps a copy of the entries it has looked up. For links, the directory would need an entry for every code, and t.co adds about 73 billion codes a year. Every follow would need a lookup in it, and a router cannot keep a copy of billions of entries. The directory would become another large, busy database. At some point, we would have to shard the directory too.
Instead of storing which shard holds each code, we can compute it from the code every time we need it. This is the same idea as a hash table, which applies a hash function to a key and takes the remainder after dividing by the number of buckets. Here, each shard plays the role of a bucket. With shards, numbered to , the router computes:
shard = hash(code) mod N
For example, with 3 shards, A, B, and C, numbered 0, 1, and 2:
| Code | hash(code) |
mod 3 |
Shard |
|---|---|---|---|
x7Kp2Qa |
17 | 2 | C |
b3Rt9Lm |
52 | 1 | B |
Qz81fWc |
30 | 0 | A |
These hash values are made up and kept small. A real hash function returns a much larger number. Similar codes give unrelated numbers.
There is no directory to store or keep up to date. Every router computes the same shard for the same code.
With Shopend’s directory, we could move one shop to its own shard. With a hash, we cannot move a code, because the formula decides its shard. That is fine for links. We handle a viral link with the cache, so we do not need to move it.