Skip to content

What is Database Sharding?

Data & Analytics, explained by the engineers who build it. Definition, how it works, use cases and common questions.

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.

Database Sharding: common questions

Something else on your mind? Ask a consultant and get a reply within one business day.

When should you shard a database?

Shard when a single primary database can no longer handle write volume or data size even after optimization, caching, read replicas and vertical scaling, or when data residency requires keeping customers' data in different regions. Most applications never reach that point, so sharding early usually adds cost without benefit.

What is a hot shard?

A hot shard receives far more traffic than the others, often because the shard key concentrates activity, such as sharding by date so all new writes hit the newest shard, or one very large tenant sharing a shard with small ones. Hot shards become bottlenecks and usually require resharding or splitting.

Does sharding improve read performance?

It can, because each shard holds less data and serves fewer queries. But queries that need data from every shard, such as global reports, become slower and more complex. If reads are the main problem, read replicas and caching are usually simpler and more effective than sharding.

Keep exploring the data & analytics glossary

Need Database Sharding in your product?

A solutions consultant replies within one business day with next steps, a rough estimate and a suggested team.