Sharding and Partitioning: Splitting Data Without Splitting Your Sanity

One database eventually runs out of one machine. Sharding buys headroom — but the shard key you pick is a near-permanent decision. Here's how to choose it.


Sharding overview: a router maps a shard key to one of three shards; a skewed key sends disproportionate traffic to one hot shard, with a comparison of good shard keys (high cardinality, even access) versus bad ones (low cardinality, monotonic)

A single database is wonderful right up until it isn’t. One machine has a ceiling — so much RAM, so many IOPS, so much write throughput — and when your data or traffic outgrows it, you have to spread across machines. Sharding is how. The mechanics are simple; the one hard decision is which key you split on, and you make it before you have the traffic that would tell you if you’re right.

Partitioning, and the two directions you can cut

Partitioning just means splitting a table into smaller pieces. There are two axes.

Vertical partitioning splits by columns — put the hot, frequently-read columns in one table and the cold, bulky ones (a big JSON blob, an avatar) in another. It’s really a schema-design move, and it keeps everything on one database.

Horizontal partitioning splits by rows — rows 1–1M here, rows 1M–2M there — each piece with the same schema but a subset of the data. When those pieces stay on one machine, it’s just partitioning (Postgres calls it declarative partitioning, and it helps with query pruning and maintenance). When the pieces live on different machines, that’s sharding, and it’s the only one of these that adds real capacity: more cores, more memory, more disk, more write throughput, in proportion to the shard count.

So: sharding = horizontal partitioning across servers. That’s the whole definition. Everything hard about it comes next.

The shard key is the entire game

A shard holds a subset of rows, and something has to decide which row goes where. That something is the shard key — a column (or set of columns) whose value maps each row to a shard. A router (a proxy, a client library, or the database’s own coordinator) computes that mapping on every read and write.

Pick this key well and the system scales almost linearly. Pick it badly and you’ve built a distributed system with the capacity of its single busiest node. Three properties separate the two:

High cardinality. The key needs many distinct values so data spreads finely. Sharding on a boolean, or on country where 60% of rows say US, gives you a few enormous partitions and can’t get finer. user_id has millions of values; status has five.

Even access. Cardinality isn’t enough — the traffic has to spread too. A key can have millions of values and still be terrible if one value is red-hot.

Query alignment. Your common queries should carry the shard key, so the router can send them to one shard. A query without the key has to ask every shard and merge — a scatter-gather whose cost grows as you add machines, which is backwards.

tenant_id and user_id usually check all three boxes. That’s not a coincidence — multi-tenant and per-user workloads are naturally partitionable because most queries are already scoped to one tenant or one user.

The two ways to map keys to shards

Once you have a key, you need a function from key to shard. There are two families, and the choice is a real trade-off.

Range partitioning assigns contiguous key ranges to shards: A–H on shard 1, I–P on shard 2, and so on. Range queries stay efficient (WHERE created BETWEEN … hits few shards), and it’s easy to reason about. The danger is skew — if keys arrive in order, one shard is always the busy one (see monotonic keys below).

Hash partitioning runs the key through a hash and assigns by the result. This spreads data evenly almost by definition, which kills hotspots from ordered inserts — but it destroys range locality, because adjacent keys land on unrelated shards, so a range scan becomes a scatter-gather.

Naive hashing (hash(key) % N) has a nasty property: change N — add or remove a shard — and almost every key remaps, forcing a near-total data reshuffle. The standard fix is consistent hashing, which moves only a small fraction of keys when the shard set changes. If you’re building hash-based sharding, you want consistent hashing underneath it; the mechanics are worth understanding on their own.

Hotspots: the failure that mocks your cluster

A hotspot is one shard drowning while the others idle. Two causes dominate.

Skewed keys. Shard by user_id and most users are fine — but the celebrity with 40 million followers, or the enterprise tenant that’s 100× everyone else, hashes to one shard, and that shard runs hot no matter how many machines you own. Mitigations exist (split the whale across sub-keys, give large tenants dedicated shards), but they’re patches on a key that didn’t anticipate the skew.

Monotonic keys. This one catches almost everyone. Shard on a timestamp or an auto-increment ID with range partitioning, and every new row has a larger key than the last, so every write lands on the newest shard. You have ten machines and a ten-way split, and nine of them are watching one do all the work. It is one of the most common sharding mistakes, and it’s why high-write systems reach for hashed keys or random prefixes on anything time-ordered.

The through-line: a hotspot converts “N machines of capacity” back into “one machine of capacity,” which is the exact thing you sharded to escape.

What you give up

Sharding is not free scaling — it trades away things a single database gave you.

Cross-shard queries get expensive. Any query that doesn’t include the shard key fans out to all shards and merges. Your reporting and search queries are usually the ones that don’t carry the key, so they get slower as the cluster grows. The common answer is a separate system for those — a search index, an analytics warehouse — fed from the shards, rather than making the operational database do fan-out work.

Cross-shard transactions get hard. A transaction touching rows on two shards can’t rely on one machine’s atomic commit. You need a distributed protocol like two-phase commit, which adds round-trips and new ways to get stuck. The whole reason a good shard key aligns with your access pattern is to keep transactions — like queries — inside a single shard, where they stay cheap and atomic.

Rebalancing is an operation, not a config change. Adding a shard means moving data while the system serves traffic, and doing it without dropping writes or reads takes real machinery. Consistent hashing limits how much moves, but “some data physically relocates between machines” is never nothing.

Referential integrity weakens. Foreign keys across shards generally aren’t enforced by the database anymore; that guarantee moves into your application, which is a step down in safety you should take on purpose.

When to shard — and when not to

Sharding is one of the highest-regret architectural decisions there is, because the shard key is so hard to change later. So exhaust the cheaper options first:

  • Vertical scaling. Modern single instances go to hundreds of cores and terabytes of RAM. “Buy a bigger box” is unglamorous and often the right call for years.
  • Read replicas. If you’re read-bound, not write-bound, replication solves it without sharding’s complexity — the writes stay on one primary.
  • Caching. Take read load off the database entirely for hot keys.
  • Move cold data out. Archive old rows to cheaper storage so the operational set stays single-node-sized.

Shard when you’re genuinely write-bound or storage-bound past a single machine’s ceiling, when the cheaper levers are spent, and — critically — when you can identify a shard key that fits your real access pattern. If you can’t name that key with confidence, you’re not ready to shard yet; you’re ready to study your queries until the key is obvious.

The rule worth remembering

Sharding buys capacity by splitting data across machines, and the shard key spends that capacity — or wastes it. Choose a key with high cardinality, even access, and alignment to your common queries, keep transactions and hot queries inside a single shard, and treat the key as close to permanent. Everything painful about a sharded system traces back to a key that didn’t match how the data is actually used.

Frequently asked questions

What is the difference between partitioning and sharding?

Partitioning is the general idea of splitting a table into smaller pieces by some rule. Sharding is horizontal partitioning where the pieces live on different machines, so it adds capacity — more CPU, RAM, disk, and write throughput — rather than just organizing data on one node. Vertical partitioning splits a table by columns; horizontal partitioning (and sharding) splits it by rows. In everyday usage 'sharding' means 'horizontal partitioning across servers.'

How do I choose a good shard key?

A good shard key has high cardinality (many distinct values so data spreads finely), even access (no single value takes a disproportionate share of traffic), and appears in your most common queries (so the router can target one shard instead of fanning out to all of them). tenant_id or user_id usually fit. Avoid low-cardinality keys like status or country, and avoid monotonically increasing keys like timestamps or auto-increment IDs, because every new write lands on the same last shard and turns your cluster into one busy node.

What is a hotspot in a sharded database?

A hotspot is a single shard receiving a disproportionate share of traffic while others sit idle, usually caused by a skewed shard key. The classic cases are a celebrity user or a giant tenant whose key hashes to one shard, and a monotonic key where all recent writes concentrate on the newest partition. Hotspots defeat the point of sharding — you paid for N machines but one of them is the bottleneck — and they are hard to fix after the fact, which is why key choice matters so much up front.

Why are cross-shard queries and transactions expensive?

Because the data no longer lives in one place. A query that doesn't include the shard key must fan out to every shard and merge the results, so its cost grows with cluster size instead of shrinking. A transaction touching rows on multiple shards needs a distributed commit protocol (like two-phase commit) to stay atomic, which adds coordination round-trips and new failure modes. Both are why the goal of a shard key is to keep the common queries single-shard.

Comments