Databases

Database Sharding Explained

What database sharding is, when you need it, how to choose a shard key, range vs hash vs directory sharding, and the operational pain points to plan for.

Server racks split into separate groups representing database shards
Illustration: Backend Architect / AI-generated.

Key takeaways

  • Sharding splits one dataset across many databases so writes and storage can scale horizontally.
  • The shard key is the most important decision: it must spread load evenly and match your queries.
  • Try indexing, caching, replicas and vertical scaling before sharding.
On this page

Sharding (horizontal partitioning) splits a large dataset across multiple database servers, each holding a subset of the rows. It’s how systems scale writes and storage beyond what a single machine can handle, and it’s one of the most consequential architectural decisions you can make.

Do you actually need sharding?

Sharding adds real complexity. Try these first:

  1. Indexes and query optimisation. Our guide to database indexing covers the basics.
  2. Caching hot reads; see caching strategies.
  3. Read replicas for read-heavy workloads.
  4. Vertical scaling: modern servers are very large.
  5. Archiving old data.

Shard when a single primary can’t handle your write volume or data size even after these steps.

Sharding strategies

StrategyHow it worksProsCons
RangeRows assigned by key ranges (A–F, G–M…)Efficient range queriesHot spots when keys cluster
HashA hash of the key picks the shardEven distributionRange queries hit every shard
DirectoryA lookup service maps keys to shardsFlexible, easy rebalancingExtra lookup; directory must be highly available
GeographicRows placed by regionData locality, complianceUneven regional load

Hash sharding with consistent hashing is a common way to add shards without moving most of the data.

Choosing a shard key

A good shard key:

  • Distributes load evenly across shards
  • Matches your most common queries, so most requests hit a single shard
  • Keeps related data together, for example sharding by customer ID so one customer’s orders live on one shard
  • Has high cardinality, with many possible values

Bad choices include timestamps (all new writes hit one shard) and low-cardinality fields like country.

The hard parts

  • Cross-shard queries and joins become expensive or impossible.
  • Transactions across shards need distributed coordination, or you redesign to avoid them.
  • Rebalancing when shards grow uneven requires moving data carefully.
  • Hot keys such as a celebrity account can overload one shard regardless of strategy.
  • Operations: backups, schema changes and monitoring multiply.

Sharding in practice

Many teams use databases or proxies with built-in sharding rather than building it themselves. Whatever you choose, design the shard key around your access patterns from the start, because changing it later is painful. Our SQL vs NoSQL guide covers which databases shard automatically.

Frequently asked questions

What’s the difference between sharding and partitioning?

Partitioning is the general idea of splitting data; sharding usually means partitioning across multiple servers.

What’s the difference between sharding and replication?

Replication copies the same data to multiple servers for availability and read scale; sharding splits different data across servers for write and storage scale. Large systems use both.

Can I shard a relational database?

Yes. It’s common, though it limits joins and transactions across shards.

Sources

  1. Martin Kleppmann — Designing Data-Intensive Applications

Every article is edited by a human and checked against our editorial policy. Spotted a mistake? Tell us.

Keep reading

Databases

SQL vs NoSQL: How to Choose

SQL vs NoSQL explained: data models, consistency, scaling, query flexibility and when to use relational, document, key-value, wide-column or graph databases.

2 min read

Databases

Database Indexing Explained

How database indexes work: B-tree and other types, composite and covering indexes, why queries ignore indexes, the write cost and checking with EXPLAIN.

4 min read