Data Sharding and Partitioning
Data Sharding and Partitioning
Definition: The practice of horizontally partitioning a large database across multiple independent database nodes (shards) to distribute traffic and storage beyond what a single machine can handle.
How It Works
- Shard key: the attribute used to decide which shard receives a given record, for example
hash(user_id) % num_shardsor a range oforder_datevalues. Every read and write needs the shard key to route efficiently. - Horizontal partitioning (sharding): splitting rows across servers, each shard holds a subset of the total rows but the full schema. This is what “sharding” usually refers to.
- Vertical partitioning: splitting columns across servers instead, for example putting rarely-read large text/blob columns on separate storage from frequently-read numeric columns. Complements sharding rather than replacing it.
- Hash-based sharding: applies a hash function to the shard key and assigns the result to a shard, giving a roughly even distribution but destroying any natural ordering, which makes range queries expensive (they must scatter-gather across every shard).
- Range-based sharding: assigns contiguous ranges of the shard key to each shard (e.g., users A-M on shard 1, N-Z on shard 2), preserving range-query efficiency but risking hotspots if the key distribution or write pattern is skewed (e.g., all new rows land on the last shard if the key is a timestamp).
- Directory-based (lookup) sharding: a separate metadata service maps each key (or key range) to a shard explicitly, giving full control over placement at the cost of an extra lookup hop and a new component that must itself be highly available.
- Geo-sharding / consistent-hash sharding: places data by region for data-residency or latency reasons, or uses a consistent-hashing ring so shard membership changes only move a fraction of the data instead of nearly all of it.
- Cross-shard queries (joins, aggregations spanning shards) require scatter-gather: fan out the query to every relevant shard, then merge results in the application or a query-routing layer, adding latency and complexity that a single-node database never has to deal with.
- Cross-shard transactions need a coordination protocol (two-phase commit, or an application-level saga) because no single database engine can natively guarantee atomicity across independent shard processes.
- Rebalancing (splitting an overloaded shard, or adding new shards as data grows) is an operational process, not a single API call: it typically means copying a range of data to a new shard, verifying it, then cutting over reads/writes with minimal downtime.
- Worked example (hash-based, 4 shards): with
shard = hash(user_id) % 4, a user withhash(user_id) = 1057lands on1057 % 4 = 1, shard 1. Going from 4 shards to 5 changes the modulus, so1057 % 5 = 2, shard 2, meaning nearly every user’s shard assignment changes at once, exactly the problem consistent hashing or a directory-based scheme avoids by only remapping the fraction of keys that need to move. - Worked example (range-based, by signup date): shard 1 holds
2023-01-01to2023-06-30, shard 2 holds2023-07-01to2023-12-31, and so on. A query for “all users who signed up in March 2023” hits exactly one shard; a query for “all users” still has to fan out to every shard, but at least date-range queries stay cheap, which is the specific trade the range scheme is bought for.
Worked Example: Capacity Planning a Shard Count
A table currently holds 4TB of data and 200,000 writes/second at peak, and a single database node is comfortable up to about 500GB and 20,000 writes/second before latency degrades. That implies a floor of 4TB / 500GB = 8 shards for storage and 200,000 / 20,000 = 10 shards for write throughput, so the binding constraint is write throughput, not storage, meaning at least 10 shards are needed today. Planning for 3x growth over the next two years pushes that to 30 shards, which is the number that actually should drive the initial shard count and shard-key design, choosing 10 “because that’s what’s needed right now” all but guarantees a disruptive resharding project well before the two years are up.
Trade-offs
- Sharding buys horizontal scalability (more machines instead of a bigger one) at the cost of losing single-node guarantees: no native cross-shard transactions, no free cross-shard joins, no single point that sees the whole dataset.
- Range-based sharding preserves efficient range scans but is more prone to hotspots; hash-based sharding balances load evenly but makes range queries expensive.
- A directory-based scheme gives the most rebalancing flexibility but adds an extra network hop and a new critical dependency; a purely algorithmic scheme (hash or consistent hash) removes that dependency but makes ad hoc rebalancing harder to control precisely.
- More shards mean more failure domains to operate (backups, monitoring, schema migrations all multiply), in exchange for smaller blast radius when one shard does fail.
- Choosing the shard key is largely a one-way door: changing it later means re-sharding the entire dataset, so the choice trades short-term convenience against long-term flexibility.
Why It Matters
- It’s the primary technique for scaling a database past the point where a single machine’s CPU, memory, disk, or IOPS becomes the bottleneck, since vertical scaling (a bigger box) has a hard ceiling and a much worse cost curve near that ceiling.
- It bounds the size of any single point of failure: losing a shard affects only the slice of data (and users) on it, not the entire dataset.
- It’s a prerequisite for multi-region or multi-tenant architectures where data locality, residency requirements, or per-tenant isolation matter as much as raw scale.
Common Pitfalls
- Choosing a poor shard key that causes hotspotting: a single shard absorbing a disproportionate share of write traffic (e.g., sharding a time-series table by a monotonically increasing ID or timestamp, so all new writes land on the newest shard).
- Picking a shard key that doesn’t match the application’s actual query patterns, forcing scatter-gather on the majority of queries instead of the minority.
- Under-provisioning for future growth so that resharding becomes necessary far sooner than planned, and resharding is disruptive precisely because it wasn’t designed in from the start.
- Ignoring cross-shard transaction needs until they show up in production, then bolting on distributed transactions or sagas as an afterthought instead of designing the data model to avoid needing them.
- Assuming sharding is the first lever to pull for scale. Indexing, caching, read replicas, and query optimization are usually cheaper and should be exhausted before taking on sharding’s operational complexity.
- Treating shard count as fixed forever. Systems that never plan a resharding path eventually hit a wall when a “final” shard count turns out not to be final.
- Sharding by a key that correlates with a business dimension nobody wants concentrated on one machine, like sharding by
signup_cohortand then running a company-wide promotion that drives a traffic spike concentrated in the newest cohort’s shard. - Not accounting for uneven data size per key even when request rate is even, a shard key that balances write traffic can still leave one shard storing far more total data if some keys (tenants, users) simply have much more historical data than others.
Comparison
| Sharding (Horizontal Partitioning) | Vertical Partitioning | Database Replication | |
|---|---|---|---|
| What’s split | Rows, across nodes | Columns, across nodes/storage | Nothing, full copies made |
| Primary goal | Write and storage scalability | Storage/IO efficiency by access pattern | Read scalability and fault tolerance |
| Cross-node joins | Expensive, scatter-gather | Cheap if co-located, else a join across stores | N/A, each replica has the full dataset |
| Adds capacity for | Both reads and writes | Mostly storage/IO efficiency | Reads only (writes still go through the leader in most designs) |
| Failure blast radius | Limited to one shard’s data/users | Limited to the split-off columns | None, every replica has everything |
| Typically combined with | Replication (each shard replicated) | Sharding or replication of the split tables | Sharding (each shard replicated) |
Sharding Strategies at a Glance
| Strategy | Distribution | Range queries | Rebalancing cost |
|---|---|---|---|
| Hash-based | Even | Expensive (scatter-gather) | Can be high unless consistent hashing is used |
| Range-based | Uneven if key is skewed | Efficient | Low, split one range into two |
| Directory-based | Fully controllable | Depends on directory design | Low, just update the mapping |
Real-World Scenario
A multi-tenant SaaS product starts on a single Postgres instance. As enterprise customers grow, a handful of large tenants start generating most of the write load, degrading performance for every tenant on the same instance. The team shards by tenant_id: each tenant’s data lives entirely on one shard, chosen via a directory service that also lets a single oversized tenant be moved to its own dedicated shard without touching anyone else’s data. Cross-tenant queries (used only by internal analytics) run through a separate scatter-gather job rather than the live request path, keeping normal application queries single-shard and fast.
Debugging Walkthrough: A Hot-Shard Incident
- Alerting fires on elevated p99 write latency for the orders service. The dashboard breaks it down by shard and shows shard 7’s CPU and disk IOPS pegged near 100%, while shards 1 through 6 and 8 through 12 sit at a comfortable 20-30%.
- First check: what’s the shard key, and could this shard’s key range explain the skew? The orders table shards by
hash(order_id) % 12, andorder_idis a monotonically increasing integer generated by a single sequence, not a UUID. - That’s the smoking gun for a subtler variant of hotspotting: even though the shard function is a hash, if
order_idincrements predictably and the hash function happens to cluster consecutive IDs onto a narrow output range for a period (or, more commonly in a real incident, the ID generator itself was recently reconfigured and started allocating from a range that happens to hash mostly onto shard 7), new writes concentrate on one shard instead of spreading evenly. - Confirming the query pattern: a query against shard 7 shows the traffic is almost entirely inserts of new orders, not reads, which points squarely at the write path and the ID-generation change, not a change in read behavior.
- Short-term mitigation: throttle or queue writes destined for shard 7 to protect it from falling over, buying time without an emergency resharding.
- Root-cause fix: replace the naive
hash(order_id) % 12with a shard function seeded from a value with better entropy (e.g., hashing a UUID or a composite key), or move to consistent hashing so future shard-count changes don’t require a full re-shard, then verify the fix by watching per-shard write-rate variance drop back to within a few percent across all 12 shards.
FAQ
Is sharding the same as partitioning? “Partitioning” is the general term; “sharding” specifically means partitioning across separate database nodes (as opposed to partitioning within a single database instance, which some systems also support).
Does sharding replace the need for replication? No, they solve different problems and are normally combined: each shard is typically also replicated for fault tolerance and read scaling, so a sharded cluster is really N shards × M replicas each.
Can you shard a database without changing the application? Rarely cleanly. Most applications need to become “shard-aware” for anything beyond simple key lookups, since joins, transactions, and unique constraints that used to be free on one node now need explicit handling.
Does every table in a sharded database need to be sharded? No. Small, rarely-written reference tables (country codes, plan tiers) are commonly replicated in full to every shard instead, avoiding cross-shard lookups for data that barely changes.
What happens to auto-incrementing primary keys under sharding? They break, since each shard would generate colliding IDs independently. Sharded systems typically use UUIDs, a centralized ID generator (e.g., Snowflake-style IDs), or a shard-prefixed key instead.
History
- Horizontal partitioning has roots in early distributed database research from the 1980s, but it became a mainstream web-scale technique in the mid-2000s as social networks and SaaS platforms outgrew single-machine relational databases.
- Google’s Bigtable (2006) and later Spanner formalized automatic range-based sharding with dynamic splitting, directly influencing HBase and other wide-column stores.
- The manual, application-managed sharding style (a fixed
hash(id) % nor a directory service in front of many MySQL/Postgres instances) was popularized by large web companies in the late 2000s before purpose-built distributed databases with native sharding (Vitess, CockroachDB, Citus) matured enough to do it automatically.
Common Interview Questions
- What makes a shard key “good” versus “bad”? A good key spreads writes evenly, matches the dominant query pattern, and rarely needs to change; a bad key concentrates traffic or forces most queries to scatter-gather.
- How do you handle a transaction that spans two shards? Either avoid the need via data modeling (co-locate related data on the same shard), or use a distributed transaction pattern like two-phase commit or a saga.
- How would you reshard a live system with zero downtime? Dual-write or use change-data-capture to copy data to the new shard layout in the background, verify consistency, then atomically cut over reads and writes.
- Why is hash-based sharding bad for range queries? Because the hash destroys the key’s natural ordering, so a range of keys is scattered randomly across all shards instead of sitting on one or two.
- What’s the difference between sharding and horizontal read replicas? Replicas hold the full dataset and scale reads; shards each hold a slice of the dataset and scale both reads and writes.
Example
Sharding a multi-tenant SaaS database by tenant_id so each enterprise customer lives on an isolated shard, letting a single noisy or oversized tenant be rebalanced onto dedicated hardware without affecting others. Instagram famously shards Postgres by using composite IDs that embed a shard identifier directly in the primary key, avoiding a separate lookup service.
Design Checklist
- Does the shard key match how the application actually queries the data, or does most traffic still need to scatter-gather across shards?
- Is there a concrete resharding plan (splitting a shard, adding shards) that doesn’t require a full rewrite of the routing logic later?
- Are cross-shard operations (joins, transactions) actually necessary, or can the data model be adjusted to co-locate related data on the same shard?
- Has sharding been chosen only after exhausting cheaper levers (indexing, caching, read replicas, query optimization)?
Related Terms
Referenced by