Data Engineering & Real-Time Event StreamsPlaybook3 min readUpdated September 2026

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.
Executive Capability Standard

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)

1. Learn:List your highest-frequency and most latency-sensitive queries and evaluate whether a candidate sharding key would let each one hit a single shard.
2. Do Manually:Set up per-shard load monitoring separate from aggregate cluster metrics so a hot shard is visible before it becomes an incident.
3. Delegate:Assign an engineer to document the reasoning behind your current sharding key and the tradeoffs it accepted.
4. Automate:Build alerting on shard load skew so an emerging hot shard triggers a notification well before it approaches capacity.
5. Buy:Bring in outside database architecture expertise before your first sharding decision if the cost of re-sharding later would be especially high for your system.

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