Adding a Shard

In the previous chapter, we designed t.co, the URL shortener X uses for every link in a post. It maps a short URL such as https://t.co/x7Kp2Qa to a long URL. The part after t.co/, here x7Kp2Qa, is the short code. We sharded t.co’s links by their short code. Following a link is the most frequent operation, and it looks up one code. So each follow goes to one shard.

For Shopend, the router finds a shop’s shard in a directory. For t.co, a directory would need an entry for every code, and t.co adds about 73 billion codes a year. So instead of looking up the shard, the router computes it from the code. With shards, numbered to , the router computes:

shard = hash(code) mod N

Here are five links on three 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
Hn4vR8e 71 2 C
pL2wJ5s 96 0 A

As before, the hash values are made up and kept small. A real hash function returns a much larger number.

A fourth shard

t.co stores about 36 TB of new links a year, so its shards fill up. At some point, we add a fourth shard, D, numbered 3. The formula becomes hash(code) mod 4:

Code hash(code) mod 3 mod 4 Moves?
x7Kp2Qa 17 C B Yes
b3Rt9Lm 52 B A Yes
Qz81fWc 30 A C Yes
Hn4vR8e 71 C D Yes
pL2wJ5s 96 A A No

Four of the five links get a different shard. Only one of them, Hn4vR8e, goes to the new shard. The other three move between shards we already had. For example, x7Kp2Qa moves from C to B. These three moves do not free up any space overall. They happen only because the formula changed.

Five links are a small sample, but the result holds in general. A link stays on its shard only when hash(code) mod 3 and hash(code) mod 4 give the same number. That happens for 3 of every 12 hash values: the ones whose remainder after dividing by 12 is 0, 1, or 2. So about 1 in 4 links stay, and about 3 in 4 move.

In general, going from shards to keeps about 1 in links on their shard. So the more shards we have, the larger the fraction of links that move.

The Google Drive variant, short links for shared documents, has the same problem. It stores fewer new links, about 9 TB a year, but we add shards to Drive too.

Moving a link means copying it to its new shard. Suppose that after three years, t.co stores about 108 TB on three shards. Going to four shards moves about 3 in 4 links. That is about 81 TB. Only about 27 TB of that ends up on the new shard.

Earlier, with Shopend, we moved one shop from one shard to another while the shop kept taking orders. We copied its data, applied the changes made during the copy, paused its writes briefly, and switched the router to the new shard. Here we have to do the same for about three quarters of all the links. Follows and new links keep arriving the whole time. Until the copy is done, every router must keep using mod 3, because the links are not on their new shards yet.

Removing a shard has the same problem. Going from 4 shards back to 3 changes the formula again, and about 3 in 4 links move again.

We want adding a shard to move only the links the new shard takes over. With four shards, that is about 1 in 4 links.