Deciding How to Shard the Database Behind Your Pipeline
Choose a database sharding key by starting from your highest-frequency, most latency-sensitive queries and checking whether each can be answered from a single shard. Sharding is expensive to reverse, so a key picked under time pressure that forces most queries to fan out across every shard is one of the costlier mistakes to walk back.
This is a guide to making that choice deliberately, since a bad sharding key chosen under time pressure is one of the more expensive mistakes to walk back later.
How do you choose a sharding key from your query patterns?
The right sharding key is determined by how you query the data, not by which field feels like the natural identifier. List your highest-frequency and most latency-sensitive queries first, and check whether each one can be answered from a single shard given a candidate key. A key that looks clean on an entity-relationship diagram but forces most real queries to fan out across every shard is the wrong choice, no matter how intuitive it looks on paper.
Watch for Hot Shards Before They Become an Incident
Even a reasonable sharding key can produce a hot shard if your data distribution is uneven, a single large customer, a popular event type, or a skewed key range that concentrates far more traffic on one shard than the others were designed to handle. Monitor per-shard load specifically, not just aggregate cluster metrics, since an aggregate view can look healthy while one shard is quietly approaching its ceiling well before anyone notices from the outside.
For example, suppose orders are sharded by customer ID because it looks like the natural identifier, and most queries are per customer. That works until one large customer's traffic overwhelms a single shard while the others sit idle. A common mistake is to treat that as a capacity problem and simply add hardware. The decision rule is to check whether the hot key is inherent to the data or a product of the sharding choice. If it is inherent, plan a split or a compound key before the shard reaches its ceiling, and record that reasoning next to the key.
Plan the Re-Sharding Path Before You Need It
Every sharding scheme eventually needs to change, whether because of growth, a new access pattern, or a hot shard that needs splitting. Decide how re-sharding would actually work, live migration with dual-write, a maintenance window, or a background backfill, while you're designing the initial scheme, not the day you discover you need it. A sharding pattern with no credible re-sharding path is a scheme that will eventually force an emergency, high-risk migration instead of a planned one.
Do you need to shard your database yet?
Sharding adds real complexity: cross-shard queries, distributed transactions or their absence, and operational overhead of managing many instances instead of one. Before committing to it, confirm that vertical scaling, read replicas, or a smarter indexing and caching strategy genuinely can't get you further, since those options are usually much cheaper to implement and reverse than a sharding scheme you later discover was premature for your actual traffic level.
Document the Sharding Key's Reasoning, Not Just the Key Itself
A sharding key documented as a fact, without the reasoning behind why it was chosen, leaves a future engineer unable to tell whether a new access pattern is compatible with the existing scheme or requires a rethink. Record which query patterns drove the decision and which tradeoffs were accepted, so a future re-sharding decision starts from an understood baseline instead of archaeology through old pull requests and Slack threads.
Test Cross-Shard Queries Explicitly, Not Just the Common Case
Even a well-chosen sharding key leaves some queries that genuinely need to touch more than one shard, such as an aggregate report or a rare lookup by a secondary attribute. Identify those queries during design, not after someone discovers in production that a report now takes ten times longer than it used to, and decide deliberately whether each one is rare enough to accept the fan-out cost or common enough to need a secondary index or a denormalized view built specifically to avoid it.
Run these cross-shard queries against production-scale data during testing, not a small local dataset. A fan-out query that's imperceptibly slow against a thousand test rows can be genuinely painful against your real shard sizes, and that gap only shows up if you test at something close to real scale, using data volumes and distribution that resemble what production will actually look like once the system has been running for a while under real, sustained usage.
Before committing to a sharding scheme, confirm that:
- The highest-frequency and most latency-sensitive queries can each be answered from a single shard under the candidate key.
- Per-shard load is monitored, not only aggregate cluster metrics, so a hot shard shows up early.
- A credible re-sharding path exists, such as live migration with dual-write, a maintenance window, or a background backfill.
- Vertical scaling, read replicas, and better indexing or caching genuinely cannot get you further.
- The reasoning behind the key and the accepted tradeoffs are written down, along with the cross-shard queries you tested.
What Good Looks Like
A sound sharding decision picks a key based on actual, documented query patterns, monitors per-shard load specifically to catch hot shards early, and has a credible, planned re-sharding path defined before it's ever needed.
Building The Capability (5-Stage Skill Ladder)
How to Get Started
Frequently Asked Questions
How do we pick the right sharding key for our database?
Start from your highest-frequency and most latency-sensitive queries, and check whether each one can be answered from a single shard given a candidate key, rather than picking whichever field looks like the natural entity identifier. A key that forces most real queries to fan out across shards is the wrong choice regardless of how clean it looks structurally.
What's a hot shard and how do we catch one before it causes an incident?
A hot shard is one receiving disproportionately more load than the others, often from a large customer or a skewed key range. Monitor per-shard load explicitly, not just aggregate cluster metrics, since an aggregate view can look fine while one specific shard is approaching its ceiling.
Do we need to plan our re-sharding strategy before we actually shard for the first time?
Yes. Every sharding scheme eventually needs to change as data grows or access patterns shift, and deciding how re-sharding will work, live migration, a maintenance window, or a background backfill, is much safer to plan upfront than to figure out under pressure once a shard is already at capacity.
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.
Blue-Green, Canary, or Rolling: Deploying Stream Processors
A decision guide to rolling, blue-green, and canary deploys for stateful stream processors, plus the rollback plan most teams never actually test.
Verifying Every Service That Talks to Your Pipeline
Which parts of zero-trust verification to build and which to buy, so every producer and consumer on a streaming pipeline proves its identity.
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.
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.