Balancing Shops Across Shards

Shopend’s performance targets apply to every shop, even while another shop is busy with a sale. Shops on the same shard still share that shard’s capacity. If one shop uses most of it, database operations for the other shops on that shard can slow down too.

Assigning shops to shards

Suppose we give each shard the same number of shops. Shopend has 10,000 shops, but they differ in size. Most are small local shops, a few hundred are regional chains, and two are national retailers.

A national retailer’s sale could overload its shard even if that shard has the same number of shops as the others. The small shops sharing it would then miss their performance targets too.

For Shopend, we give each national retailer a dedicated shard. We place smaller shops on shared shards with enough capacity for their combined workloads. This keeps a national retailer’s sale from using the database capacity assigned to smaller shops.

These assignments may need to change. A small shop can grow, and a sale can increase its traffic for a few days. When a shard no longer has enough capacity for its shops, we can move some of them to another shard or add a shard. If one shop needs more capacity than a dedicated shard has, we can divide that shop’s data further with a compound shard key, such as shopId and a product identifier. As we saw with compound keys, listing that shop’s catalog or placing an order may then need several shards.

Moving a shop

Suppose we decide to move mugshop from shard A to shard B. We have to copy its documents and update its directory entry. But if the shop keeps accepting orders on A during the copy, B will be missing the changes made after those documents were copied.

One option is a maintenance window. We stop the shop’s writes, copy its documents to B, update the directory entry, and resume on B. This is simpler, but the shop cannot accept purchases while we copy its data.

To keep the shop taking orders during most of the move, we can use the following steps:

  1. Copy the shop’s documents from A to B. The router still sends the shop’s database operations to A.
  2. Read A’s log of changes, the same log used for replication, and apply the shop’s changes on B until B is nearly caught up.
  3. Pause the shop’s writes, apply the remaining changes on B, and update the directory entry to point to B. Writes during this pause wait or are retried. Resume writes on B once routers use the new entry.
  4. Keep the old documents on A until we have checked that the move succeeded, then delete them.

A and B must never accept writes for the shop at the same time. After the switch, every router must stop using the old directory entry. Otherwise, some of the shop’s writes would still reach A while others reach B. Databases with built-in sharding can manage the copying, catching up, and switch for us.

Moving a whole shard

Moving a whole shard to a new server is a different operation. We add the new server as a replica of that shard, let it catch up, and fail over to it. The shops still belong to the same shard, so their directory entries do not change.