Cloud FinOps & Infrastructure ScalingPlaybook3 min readUpdated September 2026

Sharding Solves One Problem and Creates Five Others

Sharding gets treated as the obvious next step once a database is straining under load, split the data across multiple instances and the bottleneck disappears. It does remove that specific bottleneck. It also removes the assumption that any two rows live close enough together to query, join, or transact against in one place, and that assumption was quietly load-bearing for a lot of code that now has to change.

Shard when the alternative is genuinely worse, not because it's the default next step everyone reaches for.

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.

How do you confirm write throughput is the real bottleneck?

Sharding solves a specific problem: a single database instance can't handle the write volume anymore. It doesn't help with slow queries caused by missing indexes, doesn't help with a schema that's grown unwieldy, and doesn't help with read load, which a read replica handles with far less disruption. Profile what's actually failing before reaching for sharding. Teams that shard to fix a problem sharding doesn't address end up with all of sharding's complexity and none of the relief they expected.

How do you pick a shard key that fits your access pattern?

The shard key determines which rows land together, and a mismatched key is the single most common sharding mistake: choosing customer ID when most queries actually filter by region, or choosing a key that produces wildly uneven data distribution because a small number of values, one enormous customer, one popular region, account for a disproportionate share of the data. Trace your actual query patterns before picking a key, not the field that seems like the obvious identifier.

Cross-shard queries are the tax you'll pay indefinitely

Any query that needs data from more than one shard, a report spanning all customers, a join across entities that landed on different shards, gets meaningfully more expensive and more complex to write correctly. This isn't a one-time migration cost, it's an ongoing tax on every feature that needs cross-shard access for the life of the system. Inventory which of your current queries would become cross-shard before committing to a key, since that list tells you the real ongoing cost, not just the migration effort.

Hot shards defeat the entire purpose

A poorly chosen key can concentrate a disproportionate share of traffic onto one shard, recreating the exact bottleneck sharding was supposed to solve, just now with the added complexity of a distributed system on top. Watch per-shard load after any sharding rollout, not just aggregate throughput across all shards combined, since an average that looks healthy can hide one shard running hot while others sit nearly idle.

Rebalancing has to work before you need it, not during an outage

Data distribution shifts over time even with a well-chosen key, and eventually shards need rebalancing, moving data between them without downtime. Design and test that rebalancing mechanism before the first shard actually gets hot, not as an emergency response once it does: a rebalancing process built under pressure, during a live capacity problem, is exactly the kind of high-stakes, untested migration most likely to introduce data loss or extended downtime.

A decision rule worth writing down

Shard when write throughput on a single instance is the confirmed, measured bottleneck, when a read replica and query optimization have already been tried and aren't enough, and when the team has mapped which queries will become cross-shard and accepted that ongoing cost. Stay on a single instance, scaled vertically or with read replicas, when any of those three isn't true yet, even if the word 'sharding' has started coming up in every capacity planning conversation.

Shard only when all of these are true:

  • Write throughput on a single instance is the confirmed, measured bottleneck, not an assumption.
  • A read replica and query optimization have already been tried and aren't enough.
  • The team has mapped which queries will become cross shard and accepted that ongoing cost.
  • A rebalancing mechanism has been designed and tested before the first shard actually needs it.

A worked example: the migration that looked done and wasn't

Say a team ships a sharding migration, moves all the data, and declares it complete once the application is reading and writing against the new shards correctly. Three months later, a monthly reporting job that nobody remembered still runs a query assuming single-database access, and it starts silently returning incomplete results instead of failing loudly. Inventory every scheduled job, batch process, and analytics query against the shard migration checklist, not just the live application paths, since the ones that run rarely are exactly the ones most likely to get missed.

Executive Capability Standard

What Good Looks Like

Good sharding means a confirmed write-throughput bottleneck, a shard key chosen against real query patterns, a mapped list of cross-shard queries the team has accepted, and a tested rebalancing process built before it's needed.

Building The Capability (5-Stage Skill Ladder)

1. Learn:Profile your current database bottleneck to confirm whether it's actually write throughput, not a missing index or a read load problem sharding won't fix.
2. Do Manually:Manually trace your most common queries against a candidate shard key to estimate how many would become cross-shard.
3. Delegate:Assign an owner to design and test the rebalancing process before the first shard shows signs of running hot.
4. Automate:Automate per-shard load monitoring so an imbalance surfaces before it becomes a customer-facing performance problem.
5. Buy:Bring in outside expertise for the initial shard key design if nobody on the team has led a sharding migration before.

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

Should we shard before we actually hit a write throughput limit?

No. Sharding solves a specific, measurable problem, single-instance write throughput, and adds real ongoing complexity in exchange. Sharding preemptively, before that bottleneck is confirmed, usually means paying the cost of cross-shard queries and rebalancing for a problem you don't have yet and might solve more cheaply with vertical scaling or a read replica first.

What's the most common sharding mistake?

Choosing a shard key that doesn't match the real access pattern, most often picking an identifier that seems natural, like customer ID, without checking whether most queries actually filter or join on something else entirely. The result is a system that pays sharding's full complexity cost while still needing expensive cross-shard queries for common operations.

How do we know if our shard key is producing uneven load?

Track per-shard metrics, not just aggregate throughput across all shards combined. A healthy-looking average can hide one shard running near capacity while others sit mostly idle, and that imbalance usually traces back to a small number of unusually large or popular values concentrated on a single shard.

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