Cloud Native Patternsintermediate9 min

Sharding

When one database can't keep up, split the data across many — each holding only its own slice.

Imagine a library so popular that one front desk can't check books in and out fast enough. You could hire a faster clerk, but there's a ceiling. The smarter move is to open several desks and split the work: authors A–F at desk one, G–M at desk two, and so on. No single desk handles everything, and the whole library serves far more readers at once.

Sharding does the same thing to a database. Instead of one giant database straining under every write, you split the data across many smaller databases — each responsible for its own slice.

The problem

A single database server has hard limits: so much CPU, so much RAM, so much disk. You can buy a bigger box — vertical scaling — but that runs out fast and gets expensive fast. Adding read replicas helps you serve more reads, but every write still has to land on the one primary, and every replica has to apply that same write too, or it stops being a copy. Once your write volume or your total data size outgrows the biggest machine you can rent, replicas don't save you.

That's the wall sharding is built to break. The bottleneck isn't reads you can copy around — it's a single point that must accept every write and store every row. Try it below: traffic doubles and the database runs out of room. You pick the fix, follow it through another year of growth, then flip to the other path to compare.

How it works

You pick a shard key — a column like user_id or region — and a rule that maps each key to a shard. The rule might be a range (user_id 1–1M on shard A), a hash of the key (spreads rows evenly), or a lookup table. A router sits in front: when a write or query arrives, it inspects the shard key, figures out which shard owns that data, and sends the request straight there.

Because each shard holds only a fraction of the rows and handles only a fraction of the traffic, the system scales out simply by adding more shards. Each shard can even keep its own indexes tuned to its slice.

Step through it below. First predict where one key lands. Then a single user goes viral, and you'll see the catch: the router spreads keys, not requests. Flip to Even keys to compare with a normal minute.

Tip

The shard key is a near-permanent decision. Pick one that spreads load evenly and keeps the data you query together on the same shard, so most requests touch exactly one shard. A query that has to ask every shard at once — "all orders over $100" when you shard by customer — is the slow path you're trying to avoid; an index table can turn some of those back into direct lookups.

Check yourself

Your orders table is sharded by order date, one month per shard. Inserts are slow: one database is pinned at 100% CPU while the other eleven are nearly idle. What's the most likely cause?

Hot shards and resharding

The viral user is the loud version of a common failure. A range key keeps neighbours together, which makes range scans cheap, but a key that only grows (an auto-increment id, a timestamp) sends every new row to the last range, so one shard takes all the inserts. A hash spreads those writes evenly, at the price of scattering ranges across every shard. And no key can fix one value that's hot on its own: for that you cache it, or split it into sub-keys (42#1, 42#2, …) and gather the pieces when you read.

Growing is the other trap. With hash(key) mod N, adding a shard changes N, and that changes the home of most keys, so most of your data has to move while you're still serving traffic. Systems that expect to grow avoid this with consistent hashing, or by splitting the data up front into many small fixed partitions (say 1,024) and moving whole partitions between machines. Adding a machine then moves only the partitions it takes over.

Watch out

Sharding spreads keys, not requests. Before you commit to a key, ask what share of traffic its single busiest value gets. A celebrity account, a flagship product, a default tenant or today's date can each pin one shard at 100% while the rest sit idle — and adding shards won't help, because one key always maps to one shard.

Check yourself

You shard users by hash(user_id) mod 4 and add a fifth shard, so the rule becomes mod 5. Roughly how much of the existing data has to move?

When to use it

Reach for sharding when a single database genuinely can't hold your data or absorb your write throughput, and replicas alone aren't enough. It's the workhorse behind large-scale systems where the dataset is naturally partitionable by a key — per-user, per-tenant, or per-region data.

It's not free. Cross-shard queries and transactions get complicated, rebalancing data when a shard fills up is genuinely hard, and your application (or a router layer) now has to be shard-aware. So don't shard preemptively. Most systems do fine with a beefier box and read replicas for a long time — sharding is what you reach for once you've truly outgrown them.

Key takeaways

  • Sharding splits a single logical dataset across multiple databases (shards), each holding a distinct subset of the rows.
  • A shard key decides which shard a given row lives on; choosing it well is the whole game.
  • It scales writes and storage horizontally — something a single replicated database can't do on its own, because every replica must apply every write.
  • Good keys spread load evenly and keep related data together. Sharding spreads keys, not requests: one hot key, or a key that only grows, pins a single shard while the rest idle.
  • Resharding later is painful (with hash mod N, most rows move when N changes), so think hard about the key before you have terabytes riding on it.

Keep going