Keeping a Shop Together
We need a rule that assigns each document to a shard. A shard key is the field or combination of fields used to decide which shard a document goes to. Let’s consider two choices for Shopend as an example. We could use each document’s ID to place documents independently, or use shopId to keep a shop’s documents together.
Sharding by document ID
Suppose we use each document’s ID to assign it to a shard. The mug shop’s blue mug could be on shard A, its red mug on shard B, and a cart containing both on shard C.
To list the mug shop’s catalog, the Shopend API instance now has to read products from more than one shard and combine the results. To display the cart, it reads the cart from C and its products from A and B to get their current prices.
Placing an order involves writes too. Our purchase transaction reduces stock, records the order, and marks the cart as ordered. In this example, those changes involve several shards. A transaction on shard A covers only the data on shard A. It cannot make sure that the stock change on B and the cart update on C succeed or fail together with its own changes. We would need a transaction that coordinates the changes across shards.
Sharding by shop
Suppose we use shopId instead. All of the mug shop’s products, carts, orders, and shopper accounts go to the same shard. To list its catalog or display one of its carts, the API instance then reads everything it needs from that shard.
A cart and its purchase belong to one shop, so placing an order also stays within one shard. The transaction reads the products, reduces stock, records the order, and updates the cart there.
For Shopend, we choose shopId as the shard key because it keeps the documents these operations need together. This does not mean that we make one shard per shop. Several shops can share a shard.
Other shard keys
There are other shard keys we could use. For example, we could use the shopper identifier as the shard key. This would keep a shopper’s carts and orders together. We would still need another rule to place products, because they are shared by many shoppers.
We could create a shard for each geographic region. Such a design would group data by the region it serves. This can keep operations for a region on one shard, but the data of a shop that operates across multiple regions would need to be split across several shards.
Another option is to use an order date as the shard key. This would group orders by time, which can help reports that ask for orders from a particular period. But every new order would go to the shard holding the current period, concentrating those writes on one shard. We would also need another rule to place products and carts.
A compound key uses more than one field, such as shopId and a product identifier. This would let us spread one shop’s products across shards. However, listing that shop’s catalog or placing an order could then involve several shards.
To choose a shard key, we consider which operations need their data on the same shard and how the workload is spread across the data.