Choosing a Shard Key You Won't Have to Undo Later
A good shard key spreads data and query load evenly and appears in most of your queries, and it is the one sharding decision that is expensive to change later. Every query without that key fans out across shards, and any shard that draws more traffic than the rest becomes a hot spot the other capacity can't help.
The shard key is the one decision in this whole process that's genuinely expensive to change later, so it deserves more scrutiny than anything else in the migration.
What Makes a Shard Key Good
A good shard key distributes both data volume and query load evenly across shards, and appears in the WHERE clause of most of your actual queries, not just your largest table's primary key by habit. A key that satisfies the first condition but not the second turns routine queries into expensive cross-shard fan-outs, even if storage itself is perfectly balanced.
Look at your actual query patterns before choosing, not your data model in isolation. The right key is the one your application already asks for constantly, a tenant ID in a multi-tenant system, a user ID in a consumer app, not necessarily the most "natural" identifier on the table.
Test a candidate shard key against these points:
- It spreads data volume evenly across shards, so no single shard grows far larger than the rest.
- It spreads query load evenly too, including when one large customer or a busy time zone drives uneven traffic.
- It appears in the WHERE clause of most of your actual queries, not just as the largest table's primary key by habit.
- It matches what your application already asks for constantly, such as a tenant ID in a multi-tenant system or a user ID in a consumer app.
The Hot Shard Problem
Even a key that looks evenly distributed on paper can create a hot shard in practice if your real-world traffic isn't actually uniform, a single large customer, a viral piece of content, or a popular time zone concentrating far more activity on one shard than the others. That shard becomes your system's bottleneck even while every other shard sits mostly idle.
Model this against your actual traffic distribution, not an assumption of uniformity. If you already have one customer or entity that's an order of magnitude larger than a typical one, plan for that concentration explicitly rather than discovering it after the migration.
Cross-Shard Queries Are the Cost You Can't Avoid Entirely
Some queries will always need to touch multiple shards, a global search, an admin report spanning every tenant, and no shard key choice eliminates that entirely. What a good key choice does is minimize how often that happens for your highest-frequency, most latency-sensitive queries, pushing the expensive cross-shard cases into your lower-frequency, more tolerant ones instead.
Build your cross-shard query path deliberately, with its own performance expectations, rather than treating it as an edge case that gets patched in after the fact when someone notices an admin report timing out.
A Worked Example: Sharding a Multi-Tenant System
Say your system serves many separate customer organizations, each with its own users and data, and nearly every query already filters by which organization it belongs to. Sharding by organization ID keeps almost every query within a single shard, and any customer that outgrows its shard's capacity can, in principle, be moved to its own dedicated shard later.
The one thing to model in advance is size variance between organizations: if a handful of large customers would dominate a shard's capacity while hundreds of small ones sit nearly idle, you need a plan for isolating or splitting the largest ones before that becomes an operational fire.
Why Re-Sharding Later Is So Expensive
Changing a shard key after the fact means migrating live data between shards while keeping the system available, effectively repeating the original migration's risk with real production traffic already depending on the old layout. It's not impossible, but it's a multi-week project at minimum for most teams, not a configuration change.
That asymmetry, cheap to plan carefully up front, expensive to redo later, is the whole reason this decision deserves disproportionate attention relative to almost everything else in the sharding project.
Weigh that cost honestly against the problem sharding is solving. If a simpler fix, better indexing, a read replica, a larger instance, buys you meaningful runway, it's usually worth exhausting those options before taking on a shard key decision you'll be living with for years.
What Good Looks Like
Good sharding practice means the shard key was chosen against real query patterns and real traffic distribution, not just data volume, with a deliberate plan for the cross-shard queries that will always exist and for any known hot spots before they become an operational problem.
Building The Capability (5-Stage Skill Ladder)
How to Get Started
Frequently Asked Questions
What makes a database shard key good versus bad?
A good key distributes both data and query load evenly across shards and shows up in the WHERE clause of most of your actual queries, not just the table's primary key by default. A key that balances storage but doesn't match your real query patterns turns routine reads into expensive cross-shard fan-outs.
How do we avoid a hot shard even with a well-distributed key?
Model your real traffic distribution, not just your data distribution, before choosing. A key that looks balanced on paper can still concentrate traffic on one shard if a single large customer, viral event, or time zone drives disproportionate activity, so plan explicitly for any known concentration rather than assuming uniform load.
Can we eliminate cross-shard queries entirely with the right shard key?
No, some queries, like a global search or an admin report spanning every shard, will always need to touch more than one. The goal of a good key is minimizing how often your highest-frequency, most latency-sensitive queries need that fan-out, while accepting it for the lower-frequency cases where it's more tolerable.
Why is changing a shard key after launch so much harder than choosing it upfront?
Because re-sharding means migrating live data between shards while keeping the system available, essentially repeating the original migration's risk with real production traffic already depending on the old layout. That's a significant undertaking for most teams, which is why the initial shard key choice deserves outsized scrutiny before you commit to 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.
The PgBouncer Checklist Most Teams Skip Before Production
A pre-production checklist for PgBouncer: pool mode tradeoffs, sizing against max_connections, timeouts, failover behavior, and the double-pooling mistake.
A Production Deployment Checklist That Actually Catches Problems
A stage-by-stage deployment checklist for distributed systems, covering rollback readiness, dependency ordering, and the checks teams skip under pressure.
Verifying Devices Before They Touch Production, Not After
How to build device verification into a zero-trust rollout, what actually counts as a trust signal, and where teams stop checking too early.
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.
Why Your Redis Lock Let Two Jobs Run at Once (and How to Fix It)
A walkthrough of a real double-charge bug caused by a Redis lock's TTL expiring mid-job, and the fencing-token pattern that actually fixes it.