Model Context Protocol & Agentic ArchitecturePlaybook3 min readUpdated September 2026

How to Pick a Shard Key You Won't Regret Later

Sharding a database is one of the few architecture decisions that's genuinely hard to reverse cheaply, because the shard key you pick determines how data is distributed, and changing it later usually means moving a large share of your data while the system stays live. That asymmetry, easy to get wrong, expensive to fix, is why the decision deserves more scrutiny than most teams give it under launch pressure.

Here's what actually predicts whether a shard key holds up, drawn from where these decisions tend to go wrong once real traffic and real outliers show up.

How do you check a shard key for even distribution?

A shard key needs to spread both data volume and query load evenly across shards; a key that looks reasonable in the abstract, customer ID, for instance, can still produce badly uneven shards if your actual customer base includes a small number of accounts many times larger than the rest. Model the distribution against your real data, including your largest known outliers, before committing, rather than assuming a plausible-sounding key will behave evenly in practice.

Make sure your most common queries stay on one shard

A query that has to fan out across every shard and merge the results is dramatically more expensive than one that can be routed to a single shard directly, and a shard key chosen without checking your actual query patterns often forces exactly that fan-out on your most frequent, highest-traffic queries. Pick a key that keeps your dominant access pattern, whatever a single request or session actually needs, within one shard whenever possible.

Plan for the outlier that breaks your even-distribution assumption

Every sharding scheme eventually meets an account, tenant, or entity large enough to overwhelm a single shard on its own, even with an otherwise well-chosen key. Decide ahead of time how you'll handle that case, a dedicated shard for known large entities, or a secondary splitting strategy for a shard that outgrows its peers, rather than discovering the gap in the plan only once a specific customer's growth has already become an operational problem.

What would resharding cost before you actually need it?

If your chosen key turns out to be wrong, either from uneven growth or a shift in query patterns you didn't anticipate, resharding while the system stays live is a substantial project: it typically means dual-writing during a migration window, backfilling historical data, and cutting over reads carefully to avoid serving stale or split results. Understanding that cost ahead of time is itself a reason to over-invest in getting the key right initially rather than treating the choice as easily corrected later.

Expect a live resharding project to involve these steps:

  1. Dual-write to both the old and new layout during a migration window while the system stays live.
  2. Backfill historical data into the new layout, checking it against the source as you go.
  3. Cut over reads carefully so the application never serves stale or split results.

Build in a rebalancing mechanism from the start, even if you don't need it yet

Adding shards over time, to handle overall growth rather than a single outlier, is a different and more routine operation than fixing a bad key, but only if your system was built with rebalancing in mind from the start. Retrofitting a rebalancing mechanism onto a system that never expected to add shards is close in cost to a full resharding project, so build the capability in even during the period when you only have one or two shards and don't yet need it.

Test the failure mode where one shard is unavailable

A sharded system's behavior when a single shard is down or slow matters as much as its behavior when everything is healthy, because a shard-level outage is a narrower, more survivable failure than a full database outage only if the rest of the application actually degrades gracefully instead of failing entirely. Confirm your application handles a single unavailable shard by serving what it can and failing clearly for the affected subset, rather than treating any shard failure as a full outage across all customers.

Write the shard key's assumptions down where the next engineer will find them

The reasoning behind a shard key choice, why this key, what distribution it assumed, what the plan is for an oversized outlier, tends to live only in the memory of whoever made the decision, and that knowledge disappears the moment they move to a different team or company. Document the assumptions explicitly next to the schema itself, so a future engineer evaluating whether the scheme still holds up isn't starting from nothing.

Executive Capability Standard

What Good Looks Like

A sound sharding scheme distributes both data volume and query load evenly across shards for real traffic including known outliers, keeps common queries within a single shard, has an explicit plan for outsized entities and shard-level failures, and was built with rebalancing in mind from the start rather than retrofitted after the fact.

Building The Capability (5-Stage Skill Ladder)

1. Learn:model your actual data distribution against a candidate shard key, including your largest known accounts or tenants, before committing to it
2. Do Manually:manually trace your top few query patterns against a candidate key to confirm they'd resolve to a single shard
3. Delegate:assign a specific engineer or team ownership of the sharding scheme's health, including monitoring for uneven growth across shards
4. Automate:build automated monitoring for per-shard data volume and query load so uneven growth is caught long before it becomes a resharding emergency
5. Buy:some managed database platforms handle sharding and rebalancing as a built-in feature; evaluating one of those against a self-managed scheme is worth doing before committing engineering time to build rebalancing yourselves

How to Get Started

Frequently Asked Questions

What's the most common shard key mistake?

Picking a key based on how the data is conceptually organized rather than how the actual query load and data volume are distributed in practice. A key that looks like the natural way to think about the data can still produce badly uneven shards once real usage patterns are applied to it.

How many shards should we start with?

Fewer than you might expect, if the system supports adding more later without a full redesign. Starting with a large number of shards you don't yet need adds operational complexity for no immediate benefit; starting with a rebalancing-capable design that begins small and grows is usually the better tradeoff.

Is resharding always as expensive as this guide suggests?

It depends heavily on whether the system was designed for it from the start. A system built with rebalancing in mind can add shards relatively routinely; one that wasn't typically requires a dual-write migration and careful cutover, which is where most of the expense and risk actually comes from.

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