Sharding
In short
Sharding is a way of scaling a database by splitting its data across several servers, called shards, so each one stores and handles only part of the total.
What is database sharding?
Sharding breaks one large database into smaller pieces, called shards, that live on separate servers. Each shard holds a subset of the rows, for example users A to M on one server and N to Z on another, and all shards share the same schema. Together they behave like one logical database that can store more data and handle more traffic than any single machine.
A shard key decides where each row goes. With range-based sharding, rows are split by ranges of the key, such as dates or ID ranges; with hash-based sharding, the key is run through a hash function so rows spread evenly; and with directory-based sharding, a lookup table maps each key to its shard. The application, a routing layer, or the database itself uses the shard key to send each query to the right server.
Think of a large library that splits its collection across several buildings by the author's last name: each building is smaller and less crowded, but you have to know which one to visit. Sharding is used by large web applications and is built into or available for systems such as MongoDB, Apache Cassandra, MySQL with Vitess, and PostgreSQL with Citus.
Sharding is often confused with replication. Replication copies the same data to several servers for availability and read scaling, while sharding splits different data across servers to scale writes and storage, and large systems usually combine both by replicating each shard. Because queries and transactions that span shards are slower and harder, and a poor shard key can create hot spots where one shard gets most of the traffic, teams usually shard only after indexing, caching, and bigger hardware are no longer enough.
At a glance
Key takeaways
- Sharding splits a database's rows across multiple servers.
- A shard key determines which shard stores each row.
- Common strategies are range-based, hash-based, and directory-based sharding.
- Sharding scales writes and storage; replication copies data for availability and reads.
- Cross-shard queries and a poorly chosen shard key are the main pitfalls.
Example
// Pick a shard from the user ID, so all of a user's rows live together
import { createHash } from "node:crypto";
const shards = [dbShard0, dbShard1, dbShard2, dbShard3];
function shardFor(userId) {
const hash = createHash("md5").update(String(userId)).digest();
return shards[hash.readUInt32BE(0) % shards.length];
}
const db = shardFor(42);
const orders = await db.query("SELECT * FROM orders WHERE user_id = $1", [42]);Readers ask
What is the difference between sharding and replication?
Sharding splits data so each server holds a different part of it, which scales storage and writes. Replication copies the same data to several servers, which improves availability and read capacity. Many production systems use both.
What is the difference between sharding and partitioning?
Partitioning is the general idea of splitting a table into parts, often within a single database server. Sharding is horizontal partitioning across multiple servers, so each part runs on its own machine.
When should you shard a database?
Usually only when a single server can no longer handle the data size or write load after you have tried indexing, caching, query tuning, read replicas, and larger hardware. Sharding adds lasting complexity, so it is rarely the first scaling step.
Often compared
See also
- Database ReplicationDatabases, p. 10Database replication is the continuous copying of data from one database server to others, so several servers hold the same data for reliability and scale.
- DatabaseDatabases, p. 6A database is an organized collection of data stored on a computer, managed by software that lets applications save, search, and update it efficiently.
- CAP TheoremDatabases, p. 2The CAP theorem says that if a network failure splits a distributed database, the system must choose between consistency and availability; it can't have both.
- NoSQLDatabases, p. 29NoSQL is a family of databases that store data in models other than relational tables, such as documents, key-value pairs, wide columns, or graphs.
- CacheBackend & APIs, p. 8A cache is a fast, temporary storage layer that keeps copies of frequently used data so later requests can be served quickly without repeating slow work.
- PartitioningDatabases, p. 34Partitioning splits a large table into smaller partitions by a rule such as date ranges, so queries can skip irrelevant data and old data is easy to remove.
- Consistent HashingSoftware Architecture, p. 9Consistent hashing spreads keys across a changing set of servers so that adding or removing a server moves only a small share of the keys.
Spotted a mistake or something missing on this page?Suggest an edit