Database Sharding definition
Database sharding is a technique for scaling a database horizontally by splitting its data across multiple independent servers, called shards, each holding a subset of the rows. A shard key, such as customer ID or region, decides where each row lives. Sharding lets write volume and storage grow beyond one machine, at the cost of significant extra complexity.
How does sharding work?
Instead of one database holding every customer, a sharded system might place customers 1 to 1 million on shard A, the next million on shard B and so on, or assign them by hashing the customer ID. A routing layer, in application code, a proxy or the database itself, sends each query to the shard holding the relevant rows. Each shard is a full database with its own CPU, memory and disk, usually with replicas of its own.
Because shards share nothing, adding a shard adds write capacity, which read replicas cannot do. That is why large consumer platforms and multi-tenant SaaS products shard their busiest data once a single primary database can no longer keep up with writes.
Choosing a shard key
The shard key is the most important and hardest-to-change decision in a sharded design. A good key spreads data and traffic evenly and keeps data that is queried together on the same shard. The main strategies are:
- Hash-based: hash the key and distribute evenly; good balance, but range queries touch every shard
- Range-based: assign key ranges to shards; efficient range queries, but a risk of hot spots on recent data
- Directory-based: a lookup table maps each key to a shard; flexible, but the directory becomes critical infrastructure
- Geographic or tenant-based: keep each region's or customer's data together, which also helps with data residency
Sharding vs partitioning vs replication
Partitioning splits a table into pieces, usually within one database server, for example monthly partitions of an events table in PostgreSQL, which speeds up queries and makes old data easy to drop. Sharding is partitioning across separate servers. Replication copies the same data to several servers for read scaling and high availability. Large systems often combine all three: each shard is partitioned internally and replicated for failover.
The distinction matters when reading vendor documentation, because terms vary. MongoDB and Elasticsearch call their distribution units shards, Cassandra and DynamoDB talk about partitions, and many relational databases use partitioning for single-server splits only. What matters is whether data is spread across machines and how queries find it.
Pitfalls and alternatives to try first
Sharding makes almost everything harder: queries that span shards, joins across them, unique constraints, transactions, schema migrations, backups and rebalancing when one shard grows faster than the others. A poorly chosen key creates hot shards that defeat the whole purpose, and changing the key later means moving most of the data.
Before sharding, exhaust simpler options: query optimization and indexes, caching, read replicas, vertical scaling, archiving old data and moving analytics to a separate warehouse. When sharding is genuinely needed, consider databases that shard for you, such as Citus for PostgreSQL, Vitess for MySQL, MongoDB or distributed SQL systems.
Nexzem helps teams decide when scalability work truly requires sharding, and plans the migration in stages with dual writes, data verification, gradual traffic shifting and minimal downtime when it does, so customers never notice the move happening underneath them.