Before You Shard Your Database, Try Everything Else First
Sharding, splitting a database's data across multiple physical instances by some key, solves a specific problem: a single database instance that can no longer handle your write volume or your data size no matter how it's tuned. It also introduces a category of new problems, cross-shard queries, rebalancing, and a shard key that's expensive to change once chosen, that are worth understanding before you commit to it.
Most teams that eventually shard get more mileage than expected out of simpler fixes first: better indexing, read replicas, or splitting distinct workloads onto separate databases entirely. Sharding is what's left once those options are genuinely exhausted, not the first lever to pull.
What Sharding Actually Fixes, and What It Doesn't
Sharding fixes write throughput and total data volume that a single instance can't handle, by spreading the load across multiple instances. It doesn't fix a slow query caused by a missing index, or a database doing double duty as both your transactional store and your analytics warehouse, both of which are common root causes that look like a scaling problem but aren't actually one sharding solves.
Diagnose the actual bottleneck before reaching for sharding. A database that's slow because of unindexed queries will still be slow after sharding, just distributed across more instances, which makes the underlying problem harder to find, not easier.
Cheaper Fixes to Rule Out First
Read replicas handle a read-heavy bottleneck at a fraction of the operational complexity of sharding, since they don't require picking a shard key or rewriting your data access layer to route queries. If your bottleneck is read load, not write load, replicas are almost always the right first move.
Splitting genuinely distinct workloads onto separate databases, moving analytics queries off the same instance handling live transactional traffic, for instance, often resolves what looks like a capacity problem without touching your core data model at all. Only reach for sharding once you've confirmed the bottleneck is specifically write throughput or data volume on a single logical dataset that can't be split apart by workload.
Rule these out, in roughly this order, before you shard:
- Fix indexing first, since a slow query caused by a missing index is not a sharding problem.
- Add read replicas if your bottleneck is read load, which avoids picking a shard key or rewriting your data access layer.
- Split distinct workloads, such as transactional and reporting traffic, onto separate databases.
- Try the largest available single instance with proper indexing before concluding that one instance can't handle your write volume or data size.
Picking a Shard Key: The Decision You Can't Easily Undo
The shard key determines how data spreads across instances, and a poor choice creates hot shards, one instance handling far more traffic than the others because the key doesn't distribute evenly. A common mistake is sharding by a customer ID when a handful of large customers generate most of the traffic, concentrating load right back onto a small number of shards instead of spreading it.
Changing a shard key after the fact means re-migrating your entire dataset, which is exactly why it deserves more upfront analysis than almost any other decision in this process. Model your actual traffic distribution against candidate keys before committing, not just your data distribution, since a key that spreads data evenly can still concentrate traffic unevenly.
The Query Patterns Sharding Breaks
Any query that needs to join or aggregate across shards, a report spanning all customers, a search that isn't scoped to a single shard key, gets meaningfully harder and slower once your data is split. These queries either need to fan out across every shard and merge results in application code, or require a separate system, a read-only aggregated store, built specifically to answer them.
Inventory your actual query patterns before sharding, not just your write patterns. A team that discovers its most important reporting query needs to span every shard, after sharding is already in production, is solving a much harder problem than the one it started with.
Rebalancing: The Operational Cost Nobody Budgets For
Shards don't stay balanced forever. Traffic grows unevenly, some shards fill up faster than others, and rebalancing, moving data between shards without downtime, is genuinely difficult operational work that most teams underestimate until they're doing it for the first time under pressure.
Budget for rebalancing as an ongoing operational cost from the start, not a one time migration event. Choose a sharding scheme and tooling that supports online rebalancing if you can, since a scheme that requires downtime to rebalance will eventually force that downtime during a period of real growth, which is the worst possible time to take it.
What Good Looks Like
Good sharding decisions mean the actual bottleneck, write throughput or data volume rather than an unindexed query or a mixed workload, has been confirmed first, and the shard key was chosen against modeled traffic distribution, not just data distribution.
Building The Capability (5-Stage Skill Ladder)
How to Get Started
Frequently Asked Questions
How do we know if we actually need to shard, versus just needing a bigger instance?
Check whether the largest available instance size for your database, with proper indexing already in place, still can't handle your write volume or data size. If a bigger single instance would solve the problem, that's a much simpler fix than sharding, and worth exhausting first.
What's the most common mistake teams make picking a shard key?
Choosing a key that distributes data evenly but not traffic evenly, customer ID when a few large customers dominate traffic is the classic case. Model your actual traffic distribution against candidate keys, not just how evenly your raw data would spread across shards.
Can we avoid sharding by just using read replicas instead?
Only if your bottleneck is read load. Read replicas do nothing for write throughput or total data volume on the primary instance, since every write still has to go through it. Diagnose whether your actual constraint is reads or writes before assuming replicas will solve it.
About the numbers
This guide doesn't quote a sourced benchmark. Figures in it are estimates or general guidance, so check them against your own numbers.
Related Guides
When Your Database Actually Needs Sharding, and When It Doesn't
A decision framework for whether to shard a growing database, the cheaper fixes to rule out first, and what sharding costs you once it's live.
Deciding How to Shard the Database Behind Your Pipeline
A decision guide for picking a sharding key and pattern for the database behind a real-time pipeline, and the mistakes that force a costly re-shard.
Choosing a Sharding Key: The Criteria That Matter More Than the Technology
A decision guide for picking a database sharding key, focused on the access pattern criteria that determine whether sharding helps or hurts.
When to Shard a Vector Database (and How to Pick a Sharding Key)
A decision guide for when a growing vector database actually needs sharding, and how to choose a sharding key that doesn't wreck retrieval quality.
Sharding Solves One Problem and Creates Five Others
Sharding removes a single-database bottleneck but adds cross-shard queries, rebalancing, and hot shards. A decision guide before you commit to it.
Sharding Patterns Compared: What Actually Fits Your Data
A comparison of the common database sharding strategies and the specific tradeoffs each one makes, so you pick one before a migration forces it.