Choosing a Sharding Key: The Criteria That Matter More Than the Technology
Teams reaching for database sharding tend to spend most of their evaluation time comparing sharding technologies and very little time on the decision that actually determines success or failure: which column becomes the shard key. Get that choice wrong and no amount of infrastructure sophistication fixes the resulting cross-shard query problem. Here's how to choose it well.
Vendors Covered in this Article
Disclosure: We may earn a commission if you buy through some links on this page. It doesn't change what we recommend.
Criterion 1: does most of your query traffic already filter on this column
A shard key should match the column your application already filters by on the vast majority of its queries, whether that's a customer ID, a tenant ID, or a region. If most queries already include that column in their WHERE clause, sharding on it means most queries route to a single shard cleanly. If they don't, sharding introduces a new problem: queries that now have to fan out across every shard and merge results, which is often slower than the unsharded query it replaced.
Pull a sample of your actual production query log before deciding. A shard key chosen from how the data is modeled, rather than from how it's actually queried, is the single most common cause of a sharding project that makes performance worse.
Criterion 2: does it distribute data and load evenly
A shard key that concentrates a large share of your data or traffic onto one value, a large enterprise customer among many small ones, for example, creates a hot shard that defeats the purpose of sharding in the first place: that one shard still bears a disproportionate load, just now with the added complexity of a distributed system around it.
Check the actual distribution of your candidate key across your current data before committing. A key that looks evenly distributed in theory can turn out heavily skewed in your real dataset, and that skew is far more visible before a migration than after one.
Criterion 3: is it immutable, or close to it
A shard key that can change after a record is created, a customer moving between regions, an account changing ownership, forces a cross-shard data move every time that change happens, which is one of the more operationally painful things to build reliably. Prefer a key that's set once at creation and never needs to move.
When the natural business key genuinely does change, consider a separate, immutable identifier assigned at creation specifically for sharding purposes, decoupled from the mutable business attribute, rather than accepting cross-shard moves as a routine operation.
Criterion 4: how much does it cost to run a query that spans shards
No matter how well the key is chosen, some queries, usually reporting or admin-facing ones, will need to look across every shard. Decide upfront how those queries will run: a separate read replica aggregation layer, a periodic ETL job into an analytics store, or an accepted, occasional slow query for infrequent admin use. Deciding this before sharding, rather than discovering the need for it afterward, avoids a scramble to bolt on a cross-shard query layer under pressure once someone actually needs a report that spans customers.
Be honest about how often these queries actually run. A key that's excellent for the overwhelming majority of traffic but makes a rare, important report painfully slow is still a real cost worth planning for, even if it's not the dominant factor in the decision.
When the answer is: you don't need to shard yet
A single, well-indexed database handles more scale than most teams expect, and the operational complexity sharding adds, more moving parts, harder cross-shard debugging, migration risk, is a real and ongoing cost. Before committing to a shard key, confirm that vertical scaling, read replicas, and query optimization on your current setup genuinely can't get you further. Sharding solves a specific problem; it's not a general performance upgrade, and applying it before that specific problem exists usually creates more operational burden than it removes.
Piloting the choice before committing the whole database
Where possible, shard a single high-volume table first, using the candidate key, before applying the same key across the rest of the schema. That pilot surfaces distribution skew, cross-shard query pain, and operational friction on a smaller, more reversible scale than committing the entire database up front.
Treat the pilot's results as a real decision point, not a formality on the way to a predetermined outcome. If the pilot surfaces a problem with the chosen key, it's far cheaper to pick a different one now than after the rest of the schema has followed the same pattern.
Before committing to a shard key, confirm the following:
- Most of your production queries already filter on the column, so they route to a single shard instead of fanning out.
- Data and load distribute evenly, with no single value, such as one large customer, concentrating traffic on a hot shard.
- The value is set once at creation and rarely or never changes, so records do not need to move between shards.
- You have decided how cross-shard reporting and admin queries will run, whether through a replica layer, an ETL job, or accepted slow queries.
- Vertical scaling, read replicas, and query optimization genuinely cannot take your current setup any further.
What Good Looks Like
Good sharding practice means the shard key matches real query patterns, distributes load evenly, stays immutable, and has a defined plan for the cross-shard queries that will inevitably still occur.
Building The Capability (5-Stage Skill Ladder)
How to Get Started
Disclosure: We may earn a commission if you buy through some links on this page. It doesn't change what we recommend.
Frequently Asked Questions
Can we change our shard key later if we chose wrong?
Yes, but changing a shard key is a significant migration that moves data across a new boundary while the system stays live. It benefits from the same discipline as a zero-downtime schema migration: small, reversible steps instead of a single cutover, with rollback defined at each step.
How many shards should we start with?
Fewer than you might expect. Starting with a small number of shards, each still handling a meaningful load, is easier to operate and debug than over-splitting from the start. You can add shards later; consolidating an oversplit setup is harder.
Does sharding replace the need for indexing and query optimization?
No. A poorly indexed query is still poorly indexed after sharding, just now running that way on a smaller slice of data per shard. Fix indexing and query patterns first; sharding addresses scale that a single well-tuned database genuinely can't handle, not queries that were never optimized to begin with.
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.
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.
Before You Shard Your Database, Try Everything Else First
Sharding solves a real scaling problem and creates several new ones. Here is what to rule out first, and how to pick a shard key if you do need it.
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.
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.