The short answer
Quick answer: Sharding splits one large database into smaller pieces called shards, each stored on a different server. Every row belongs to exactly one shard, chosen by a shard key such as the customer ID. This lets a database grow beyond the storage and write capacity of a single machine: add more shards to hold more data and handle more writes. The cost is complexity. Queries that span shards are slow, transactions across shards are hard, and choosing a poor shard key causes problems that are painful to undo. Shard only when simpler options have run out.
Why sharding exists
A single database server has limits:
- Storage. The data no longer fits, or backups and restores take too long.
- Write throughput. One machine can only process so many writes per second.
- Memory. Indexes no longer fit in RAM, so queries hit disk.
Replication helps with reads, but every replica still holds all the data and must apply every write. To scale writes and total size, the data itself must be divided.
Terminology: partitioning means splitting a table into parts. When the parts live on the same server, it is just partitioning. When they are spread across servers, it is sharding (horizontal partitioning).
Choosing how to split
Range-based
Assign contiguous key ranges to shards: customers A to F on shard 1, G to M on shard 2, and so on.
- Good: range queries are efficient, since neighbouring keys live together.
- Bad: uneven distribution. If keys are timestamps or sequential IDs, all new writes hit the last shard, creating a hot spot.
Hash-based
Compute a hash of the key and use it to pick a shard.
- Good: spreads data and load evenly.
- Bad: range queries must ask every shard, because neighbouring keys are scattered.
A naive hash(key) % number_of_shards has a serious flaw: change the number of shards and almost every key moves. Consistent hashing solves this by moving only a small share of keys when shards are added or removed. Many systems instead hash into a large fixed number of virtual buckets and assign buckets to servers, so rebalancing means moving whole buckets.
Directory-based
Keep a lookup table mapping each key (or tenant) to a shard.
- Good: fully flexible. A large customer can be moved to its own shard.
- Bad: the directory is an extra component that must be fast and highly available.
Geographic
Place data by region: European users on European shards. This reduces latency and can help meet data residency rules.
| Strategy | Even distribution | Range queries | Flexibility |
|---|---|---|---|
| Range | Poor without care | Excellent | Medium |
| Hash | Excellent | Poor | Low |
| Directory | Up to you | Depends | High |
The shard key is the big decision
A good shard key has:
- High cardinality: many distinct values, so data can be divided finely.
- Even distribution: no value owns a huge share of the data or traffic.
- Query alignment: most queries include the key, so they can go to a single shard.
- Stability: it does not change for a given row.
Typical choices are tenant_id for a multi-tenant product or user_id for a consumer app, so that everything for one customer or user lives together.
Poor choices and why:
| Key | Problem |
|---|---|
| Creation timestamp | All new writes go to one shard |
| Country | A few countries hold most of the users |
| Status or boolean | Too few distinct values |
| A value that changes | Rows would have to move between shards |
Even a good key can have a hot key: one celebrity account or one enormous customer that overloads its shard. Solutions include isolating that tenant on dedicated hardware or splitting its data further.
What becomes hard
Cross-shard queries
A query that includes the shard key goes to one shard. A query that does not must be sent to all shards and the results merged (scatter-gather). It is as slow as the slowest shard and uses resources on every one. Aggregations, sorting and pagination across shards are all more work.
Joins
Joins between tables on different shards are expensive or unsupported. The usual answer is to co-locate related tables on the same key, so a customer's orders and order items share a shard, and to denormalise where needed.
Transactions
A transaction touching two shards needs a distributed protocol such as two-phase commit, which is slower and has awkward failure modes. Good designs keep each transaction within one shard. See what ACID really means.
Unique constraints and IDs
A database cannot cheaply enforce uniqueness across shards, and auto-increment counters collide. Sharded systems use globally unique IDs such as UUIDs or time-ordered "snowflake" IDs.
Rebalancing
As data grows, shards must be split or data moved, while the system stays online. This involves copying data, keeping the copy in sync, and switching traffic over without losing writes.
Operations
There are now many databases to back up, monitor, upgrade and migrate schemas on. Each shard also needs its own replicas for availability.
How routing works
Something must know which shard holds which key:
- In the application. The code computes the shard and connects to it. Simple but spreads the logic everywhere.
- In a proxy or router. Applications connect to a middle layer that forwards queries. Vitess does this for MySQL, and MongoDB's
mongosrouter does it for sharded clusters, as described in the MongoDB sharding documentation. - In the database itself. Distributed databases such as Cassandra, DynamoDB, CockroachDB and Spanner partition and rebalance data automatically.
Try these first
Sharding is expensive to adopt and hard to reverse. Before you do it:
- Optimise queries and indexes. See how database indexes work.
- Scale up. Modern servers are very large.
- Add read replicas and a cache.
- Archive or delete old data.
- Partition tables within one server.
- Split by function: move a separate feature's tables to their own database.
If you do need horizontal scale, consider a database that shards natively before building your own layer. That decision is part of the wider SQL vs NoSQL discussion.
Frequently asked questions
What is the difference between sharding and partitioning?
Partitioning is splitting a table into parts. Sharding is partitioning across multiple servers.
What is a shard key?
The column (or columns) whose value decides which shard a row is stored on.
Can I change the shard key later?
It is possible but usually means migrating all the data. Treat it as a near-permanent choice.
When should I shard my database?
When a single primary can no longer handle your write volume or data size, and you have already exhausted indexing, hardware upgrades, replicas, caching and archiving.
Conclusion
Sharding trades simplicity for scale. It removes the single-machine ceiling by dividing data across servers, and in exchange every query, transaction and operational task has to think about where data lives. Pick a shard key that matches how your data is accessed, keep related data together, and delay sharding until you truly need it.
Related articles
- How Database Replication Works
- How Consistent Hashing Distributes Data Across Servers
- SQL vs NoSQL: When to Use Which (Beyond the Hype)
- Why Big Tech Companies Build Their Own Databases
