What Is a Shard Key and How Do You Choose One?
A shard key is the field, or set of fields, whose value decides which shard stores a row. A shard is one of the servers that together hold a database too large or too busy for one machine. A good shard key spreads reads and writes evenly across shards. It also lets the most common queries be answered by one shard. A bad one sends most traffic to one server, which is the same as not sharding at all.
Why a database is split into shards in the first place is explained in What is database sharding?. This page is about choosing the key.
Property 1: high cardinality
Cardinality is the number of distinct values a field can take. A shard key needs many more distinct values than there are shards. The reason is that every row with the same key value must live on the same shard. A boolean field has two values, so it can fill at most two shards. A country field has about 200 values, but most rows would share a few of them. A user_id in a system with 50 million users has 50 million values and can be divided any way the cluster needs.
Property 2: even frequency
Frequency is how often each value appears. High cardinality is not enough if one value dominates. Take a customer_id key in a marketplace where one retailer produces 30 percent of all orders. That key puts 30 percent of the data and traffic on one shard. Look at the distribution, not only the count of distinct values. If a few values are far heavier than the rest, add a second field to the key so that the heavy value is split.
Property 3: non-monotonic
A monotonic key is one whose new values always increase, such as a timestamp or an auto-increment id. Under range sharding, every new write has the highest value so far. It therefore goes to the last range, on one shard. That shard takes 100 percent of the inserts while the others sit idle, which is the classic hot shard.
Compare user_id with created_at for an orders table. Sharding on user_id sends each new order to the shard that owns that user. Inserts spread across all shards, and "show me my orders" reads one shard. Sharding on created_at sends every new order to the newest shard and spreads one user's orders across all of them. If you need time in the key, hash it or put it second, behind a field that spreads.
Hash versus range sharding
Hash sharding applies a hash function to the key value and uses the result to pick the shard. The hash scatters even a monotonic key, so writes spread evenly, and a lookup by exact key hits one shard. The cost is the range query, such as all orders between two dates. It must be sent to every shard, because neighboring values are on different servers.
Range sharding assigns each shard a contiguous range of key values, for example user_id 1 to 1,000,000 on shard 1. A query over a range of keys reads one or two shards, and sorted scans are cheap. The cost is the hot shard from a monotonic key. Ranges must also be split and moved as data grows. Choose hash when the workload is point lookups on a key. Choose range when range scans on the key are common and the key is not monotonic.
Compound keys: DynamoDB and MongoDB
A compound shard key uses two or more fields. Amazon DynamoDB makes this explicit with a partition key and an optional sort key. The partition key is hashed to choose the partition. Items with the same partition key are stored together, ordered by the sort key. A table keyed on (user_id, order_date) puts all of one user's orders on one partition, sorted by date. "This user's orders last month" is then one range read on one partition. Each partition serves at most 3,000 read units and 1,000 write units per second, so a busier partition key is throttled.
MongoDB lets you pick a ranged or a hashed shard key, and the key must be supported by an index on the collection. A query that includes the shard key is routed to one shard. A query that does not is broadcast to all shards and merged. MongoDB's own guidance repeats the three properties above and adds a fourth. Pick a key that appears in your most frequent queries, so that they are routed rather than broadcast.
What goes wrong
Hot shards. A hot shard is one that receives far more traffic than the others. The cause is a monotonic key, a dominant value, or one very popular entity. A celebrity account with 50 million followers keyed on account_id puts every read of that account on one shard. Fixes include adding a random suffix to the key for that entity, caching the hot rows, or splitting the heavy value across several shards.
Scatter-gather queries. A scatter-gather query is one that cannot be routed to one shard. The router sends it to every shard and merges the results. With 20 shards, one such query costs 20 database calls and waits for the slowest. A design that runs its most common query as scatter-gather has chosen the wrong key. The usual repair is a secondary index or a second copy of the data keyed differently.
Resharding. Changing the shard key later means every row is rehashed and most rows move to a different server. Recent MongoDB versions can reshard a collection in place. The operation still rewrites the whole collection and needs free disk space equal to its size. In a hand-built setup on MySQL or PostgreSQL, the migration is a dual-write period, a backfill, and a cutover. That is often weeks of work, which is why the key is chosen carefully at the start. The full list of costs is in What are the disadvantages of sharding?. The vocabulary difference is in What is the difference between Sharding and Partitioning?.
How to Prepare
- State the three properties and give the user_id versus created_at example without notes. That is the most common follow-up question.
- Practice on three tables (users, orders, messages). Name a key for each, the query it makes cheap, and the query it makes expensive.
- Walk through a hot shard fix out loud, from detection to the random suffix or cache that repairs it, following How to Handle Database Sharding in a System Design Interview.
- Study the data partitioning chapters in Grokking System Design Fundamentals, which cover hash, range, and directory-based sharding with diagrams.
- Choose a shard key inside a full design with Grokking the System Design Interview, where the URL shortener, the chat system, and the feed each depend on that choice.

GET YOUR FREE
Coding Questions Catalog

$99

$197

$72