AlgoMaster Logo

Sharding

High Priority15 min readUpdated September 25, 2026
AI Mock Interview

Practice this topic in a realistic system design interview

Listen to this chapter
Unlock Audio

Premium Video

This video is available to premium subscribers only

Unlock Full Access

Suppose an application runs on a single relational database. At first, this works perfectly well. The database stores users, orders, payments, and everything else the application needs.

As traffic grows, we can usually scale the database vertically by adding more CPU, more memory, faster storage, or moving to a larger machine. We can also add read replicas to take on more of the read traffic.

But eventually, one database server reaches its limit:

  • The dataset may become too large for one machine.
  • Write traffic may become too high for one primary to absorb.
  • Storage requirements keep growing.
  • Each upgrade to a larger machine costs more than the last one.

At that point, instead of keeping all the data on one database server, we may need to divide it across multiple servers. That is the basic idea behind sharding.

This chapter covers what sharding is, why systems use it, and the trade-offs that come with splitting data across multiple database servers.

1. What is Database Sharding?

Database sharding means splitting a large dataset across multiple independent database instances. Each database stores only a portion of the total data, and each of these databases is called a shard.

For example, suppose we have 100 million users. Instead of storing all of them in a single database, we could distribute them across four shards of 25 million users each.

Together, those four shards still contain the complete user dataset. But no single database has to store or process all of it. Each shard handles roughly a quarter of the storage and a quarter of the queries.

This is a form of horizontal scaling. Instead of making one database server more powerful, we add more database servers and distribute the data across them.

2. Sharding vs Replication

Sharding is sometimes confused with replication, but they solve different problems.

  • With replication, multiple database servers store copies of the same data.
  • With sharding, different servers store different subsets of the data.

So replication copies data, and sharding divides it.

ReplicationSharding
What each server storesA full copy of the dataA portion of the data
100 million usersEvery server holds all 100MEach server holds a slice
Main benefitAvailability and read scalingScaling storage and write capacity

In practice, large systems often use both. The data is divided into shards, and each shard has its own replicas for availability and read scaling.

Each primary holds only its own shard's data, and its replicas copy that shard and nothing else.

3. The Shard Key

The most important design decision in a sharded database is the shard key. The shard key determines which shard stores a particular row.

Suppose we shard a users table by user_id. When a new user is created, the system looks at that user_id and uses it to decide which shard should store the record. Later, when we want to fetch that user, the same shard key lets us route the query to the correct database:

Because the query includes user_id, the system knows which shard holds user 123 and does not need to ask the others.

A good shard key has two properties:

  1. It spreads data and traffic evenly. Every shard should hold a similar amount of data and receive a similar share of requests.
  2. It makes common queries easy to route. The queries the application runs most often should include the shard key, so each one can go to a single shard.

Choosing the wrong shard key can create serious scaling problems, even if the system technically has many shards. One shard may end up with most of the traffic, or most queries may have to visit every shard.

Once we have a shard key, we still need a rule that maps each key value to a shard. There are three common strategies.

4. Range-Based Sharding

With range-based sharding, we divide the key space into ranges and assign each range to a shard.

For example, users 1 through 1,000,000 go to Shard 1, users 1,000,001 through 2,000,000 go to Shard 2, and so on.

ShardKey Range
Shard 1user_id 1 to 1,000,000
Shard 2user_id 1,000,001 to 2,000,000
Shard 3user_id 2,000,001 to 3,000,000

This approach is simple and easy to understand. It also works well for range queries, because related values are stored close together. A query for users 1,200,000 through 1,300,000 only needs to visit Shard 2.

But it has an important weakness: traffic may not be evenly distributed.

Suppose user IDs increase over time. Every newly created user gets the next ID, so all new signups are written to the shard that holds the newest range. New users also tend to be the most active ones. That shard can end up receiving far more traffic than the older shards, which becomes a hot shard.

So range sharding preserves locality, but it can also create uneven load.

5. Hash-Based Sharding

Hash-based sharding takes a different approach. Instead of assigning ranges directly, we apply a hash function to the shard key and map the result to one of the shards.

Conceptually, the system does something like this:

If the result is 0, the row goes to Shard 0. If the result is 1, it goes to Shard 1, and so on.

Hashing tends to distribute keys more evenly than simple ranges. New users no longer pile onto one shard, because consecutive IDs hash to unrelated values. That helps spread both storage and traffic across shards.

But we lose locality. Users with nearby IDs may end up on completely different shards. A range query such as "users 1,200,000 through 1,300,000" now has to visit every shard, because the matching users are scattered across all of them.

So hash-based sharding usually improves distribution at the cost of locality.

6. Directory-Based Sharding

Range and hash sharding both calculate the shard from the key. Directory-based sharding looks it up instead.

The system maintains a lookup table that records which shard holds each user, customer, or tenant. For example:

CustomerShard
Customer AShard 3
Customer BShard 7
Customer CShard 2

This gives us much more flexibility. If Customer A grows large enough to need its own shard, we can move its data and update one entry in the directory. The shard key itself does not change, and no other customer is affected.

But now we have another piece of infrastructure to maintain. Many database requests depend on the directory, so it must be fast and highly available. If it is slow, every request is slow. If it is down or returns the wrong shard, routing breaks.

So directory-based sharding trades simplicity for flexibility. The table below compares the three strategies.

Range-basedHash-basedDirectory-based
How the shard is chosenKey falls in a rangeHash of the keyLookup table
Even distributionCan be unevenUsually evenControlled manually
Range queriesEfficientVisit every shardDepends on the mapping
Moving data between shardsSplit or move rangesHard with plain moduloUpdate the directory entry
Extra infrastructureNoneNoneA fast, highly available lookup service

You can try these strategies in the simulation below and see how each one places keys across shards.

Loading simulation...

7. Query Routing

Once the data is sharded, the application needs to know where to send each query. This job belongs to a routing layer. It can live in application code, a data-access library, a proxy, or the database system itself.

Suppose we want to fetch user 123. If user_id is the shard key, the routing layer calculates which shard owns that user and sends the query directly there. Only one database handles the request, which is the ideal case.

Now consider a different query:

signup_date is not part of the shard key, so the system has no way to know which shards contain the matching users. It has to send the query to every shard and combine the results. This is called a scatter-gather query.

Scatter-gather queries get more expensive as the number of shards grows. With 4 shards, the query does 4 times the work of a single-shard lookup. With 100 shards, it does 100 times the work, and the response is only as fast as the slowest shard.

This is why shard key design should be based on the access patterns of the application, not just the shape of the data.

8. Cross-Shard Queries

Scatter-gather is the simplest form of a broader problem. Sharding gets more complicated whenever a query needs data from multiple shards.

Suppose users are sharded by user_id, and we want the total number of active users across the entire system. Each shard only knows about its own users. So the application has to query every shard for a partial count and then add the partial counts together.

Sorting and pagination are harder still. Imagine asking for the 100 most recent orders globally:

Each shard has its own set of recent orders, and any of them could hold some of the global top 100. So the system asks each shard for its own 100 most recent orders, merges those lists, sorts them, and keeps the top 100.

Fetching page 2 is worse. The system cannot ask each shard for "orders 101 to 200", because it does not know how many of the global orders 101 to 200 each shard holds. It usually has to fetch the top 200 from every shard and discard most of them.

The more cross-shard work a query requires, the more complexity and latency it adds.

9. Cross-Shard Transactions

Transactions are another major challenge.

Suppose a money transfer moves funds between two accounts. If both accounts live on the same shard, the database can handle the transaction normally, and both updates either commit together or not at all.

But what if one account is on Shard 2 and the other is on Shard 6? Now a single logical transaction spans two independent databases, and neither one can guarantee on its own that both updates succeed together.

Coordinating that operation is much harder. The system may need a distributed transaction protocol, where the shards coordinate a single commit, or an application-level pattern such as a Saga, where each step has a compensating action that undoes it if a later step fails. Both approaches add latency and introduce new failure cases.

Because of this, well-designed shard keys try to keep data that is frequently updated together on the same shard. This is often called data locality.

For example, if orders and order_items are both sharded by user_id, creating an order and its items is a single-shard transaction. This usually means storing user_id on order_items too, even though it could be looked up through the order, so both tables use the same shard key.

10. Hot Shards

Even if data is evenly distributed across shards, traffic may still be unbalanced.

Suppose we shard a social media system by user_id. Most users receive very little traffic. But a few extremely popular accounts receive millions of requests. If one of those accounts lives on a particular shard, that shard may become overloaded while the others remain mostly idle.

This is called a hot shard. The whole system can become limited by that one overloaded shard, even when the other shards have plenty of spare capacity.

Hot shards can also come from a poor shard key choice. For example, sharding by region can put most of the traffic on one shard if the majority of users come from the same country. Sharding by a field with only a few possible values, such as status or plan_type, has the same problem.

So good sharding is not only about distributing the amount of data. The workload needs to be distributed too.

11. Resharding

Eventually, the system may need more shards.

Suppose we started with four shards, and the application has grown enough to need eight. We cannot simply add four empty databases and expect the load to spread out on its own. Some of the existing data has to move from the original shards to the new ones. This process is called resharding.

Resharding is difficult because the application is usually still serving live traffic while the data moves. During the migration, a row may exist in both its old and new location. The system has to route reads and writes correctly, avoid losing writes that arrive mid-copy, and avoid returning inconsistent results.

Some databases provide built-in support for automatic rebalancing. In other systems, the application team has to manage the migration itself.

The mapping strategy also affects how much data has to move. With plain modulo hashing, changing the number of shards changes the result for most keys. Going from 4 shards to 5, for example, keeps only about 1 in 5 keys on the same shard, so roughly 80% of the data moves, even though the new shard only needs about 20% of it.

Hash-based systems sometimes use techniques such as consistent hashing to reduce how much data needs to move when shards are added or removed. With consistent hashing, adding a shard moves roughly only the share of data the new shard should own.

12. Sharding and Failure Isolation

Sharding can also improve failure isolation.

Suppose we have ten independent shards and one of them fails. Only the users stored on that shard are affected. The other nine shards keep serving traffic, so about 90% of users see no problem at all. That is much better than having the entire dataset depend on one database server.

But there is an important caveat. Each shard still needs its own availability strategy. If a shard lives on only one machine and that machine fails, its data becomes unavailable until the machine recovers.

So production systems often combine sharding with replication. Each shard has a primary and one or more replicas, and if a shard's primary fails, one of its replicas can be promoted to take over.

Here, Shard 2's primary has failed and its replica has taken over. Shards 1 and 3 are not affected at all.

13. Operational Complexity

Sharding provides scalability, but it adds significant operational complexity. Instead of managing one database, you may be managing dozens or hundreds.

  • Schema migrations need to run across every shard, and a migration that fails halfway leaves some shards on the new schema and some on the old one.
  • Monitoring needs to track each shard independently, since one slow or full shard can hide inside healthy averages.
  • Backups become more complicated, because each shard is backed up separately and restoring a consistent point in time across all of them is harder.
  • Rebalancing data becomes a real operational task rather than a one-time setup step.
  • Debugging gets harder, because two users making the same request may be talking to completely different databases.

If different shards accidentally end up with different schema versions or configuration, subtle problems appear that only affect some users, which makes them hard to reproduce.

So sharding is usually not something to add early just because the system is expected to grow. It is a technique to introduce when simpler scaling approaches are no longer enough.

14. When Should You Shard?

Before sharding a database, it is worth exhausting simpler options:

  1. Optimize slow queries.
  2. Add indexes.
  3. Introduce caching.
  4. Scale the database vertically.
  5. Move read traffic to replicas.
  6. Archive old data.
  7. Partition large tables inside the database.

Each of these is far cheaper to adopt and operate than sharding, and together they can carry a database a long way.

Sharding becomes useful when one database can no longer handle the dataset or the write workload efficiently, and the application needs to scale horizontally.

You should also make sure the application has a natural shard key. If most operations can be routed to a single shard, sharding can work very well. A multi-tenant SaaS application, where almost every request belongs to one customer, is a good example. If almost every request needs to join data across many shards, the design may create more problems than it solves.

Summary

Sharding splits a large dataset across multiple independent databases called shards, so no single server has to store or process all of the data. It is a form of horizontal scaling, and it is different from replication: replication copies data, while sharding divides it. Large systems often use both, with each shard having its own replicas.

The shard key decides where each row lives. A good one spreads data and traffic evenly and lets common queries go to a single shard. Range-based sharding keeps related keys together but can create hot shards. Hash-based sharding spreads keys evenly but loses locality. Directory-based sharding offers the most flexibility, at the cost of a lookup service that must always be available.

Queries that include the shard key route to one shard. Queries that do not become scatter-gather queries, and aggregations, global sorting, and cross-shard transactions all add latency and complexity. Hot shards, resharding, and the operational load of running many databases are ongoing costs.

Sharding improves failure isolation, but each shard still needs replication for availability. Use it when simpler options are exhausted and the application has a natural shard key.

Quiz

Sharding Quiz

10 quizzes