Database Sharding
Splitting one huge dataset across multiple databases, so no single database has to hold all of it.
What is it?
Replication solves reliability and read scaling by copying the same data multiple times. But eventually a dataset can grow so large that even writes — or just the storage itself — outgrow what a single database server can hold, no matter how good its hardware is. Sharding solves this differently: instead of copying everything, it splits the data itself into pieces (shards), each stored on a separate database, based on some rule — commonly a shard key, like a user id.
Explain like I'm 10
It's like a massive filing system split across multiple filing cabinets by last name — A-M in one cabinet, N-Z in another. Neither cabinet holds everything, but together they hold it all, and each one only has to be big enough for its own share.
Examples
Routing to the right shard by user id
function getShardForUser(userId) {
const shardCount = 4;
return userId % shardCount; // simple hash-based routing
}
async function getUser(userId) {
const shard = getShardForUser(userId);
return databases[shard].query("SELECT * FROM users WHERE id = ?", [userId]);
}How it works
A shard key determines which shard a given piece of data belongs to — often computed with a hash function so data spreads out roughly evenly. Every read or write for a given piece of data goes to exactly the one shard responsible for it, so each individual database only ever holds and processes a fraction of the total dataset.
Why does it exist?
Some datasets are simply too large — in size or in write traffic — for any single database server to handle, no matter how powerful. Sharding is the way to scale a database horizontally (more machines) rather than vertically (a bigger machine), the same underlying idea as scalability generally.
When to use it
Reach for sharding once a dataset's size or write throughput has genuinely outgrown what a single (even well-replicated) database can handle — this is usually a late-stage scaling decision, not an early one.
When not to use it
Don't shard prematurely — it adds real complexity (queries that need data from multiple shards become much harder, and re-sharding later is painful) for a scale most applications never actually reach. Exhaust simpler options first: a bigger server, replication, and caching.
Common mistakes
Choosing a shard key that leads to uneven distribution — e.g. sharding by signup date when most users signed up recently, overloading one shard.
Writing queries that need to join or aggregate data across multiple shards, which sharding makes much more expensive or awkward.
Sharding before it's actually necessary, taking on real operational complexity for a scale that hasn't been reached yet.
Practice exercises
- Easy:
Explain, in your own words, the difference between sharding and replication.
- Medium:
Describe a bad shard key choice for a table of orders, and why it would cause uneven load.
- Hard:
Explain why a query that needs to count all rows across every shard is more expensive than the same query on an unsharded database.
Interview questions
What is database sharding?
Splitting a dataset into pieces, each stored on a separate database, based on a shard key, so no single database has to hold all the data.
What's the difference between sharding and replication?
Replication copies the same data to multiple databases; sharding splits different pieces of data across multiple databases — they solve different scaling problems and are often used together.
What makes a good shard key?
One that distributes data (and load) roughly evenly across all shards, and that most queries can use directly to route to the right shard without needing data from others.
What's the difference between range-based and hash-based sharding?
Range-based sharding assigns contiguous ranges of the key (e.g. user ids 1–1000 on shard A, 1001–2000 on shard B) to each shard, which keeps related/sequential data together and makes range queries efficient. Hash-based sharding runs the key through a hash function to pick a shard, spreading data roughly evenly but scattering anything that was sequential or related across unrelated shards.
What's a downside of range-based sharding regarding hotspots — for example, sharding orders by order id or timestamp?
Whatever range currently holds the newest/most-active data (like the highest order ids, or today's date range) receives essentially all new writes and recent reads, while older ranges on other shards sit comparatively idle — defeating the purpose of spreading load evenly.
Scenario: a table is sharded by `userId % 4`. What happens when you need to add a 5th shard?
Almost every user's userId % 5 now maps to a different shard than userId % 4 did, so nearly the entire dataset needs to be physically moved to match the new formula — a naive modulo scheme makes adding or removing shards enormously disruptive rather than incremental.
What is consistent hashing, and how does it solve the resharding problem above?
Consistent hashing maps both shard keys and shard nodes onto positions on the same conceptual ring, and each key belongs to the nearest node going around it. Adding or removing a node only reassigns the keys nearest to that one node, rather than remapping almost the entire keyspace the way a plain modulo scheme does.
What is a "hot shard," and give a concrete example of a shard key that would create one?
A hot shard is one that receives disproportionately more traffic or data than the others, becoming a bottleneck even though the system is nominally "sharded." Example: sharding by country when 80% of your users are in a single country — that one shard carries most of the real load no matter how many shards exist.
What is "directory-based" sharding, and what's its tradeoff versus computed (hash or range) sharding?
A directory-based scheme keeps an explicit lookup table mapping each key (or key range) to its shard, rather than computing the mapping with a formula. It's more flexible — individual keys can be moved or rebalanced independently — but that lookup table itself becomes a critical, potentially bottlenecked piece of shared infrastructure that computed schemes avoid.
Why is choosing a shard key one of the hardest, most consequential decisions in a sharded system?
It simultaneously determines how evenly load spreads across shards, whether your most common queries can be routed to a single shard directly, and how expensive resharding will eventually be — get it wrong and you can face hot shards, constant cross-shard queries, or a painful migration, often all discovered only once real production traffic exposes it.
Why is it hard to change a shard key after a system is already in production?
Changing the shard key means recomputing which shard every single row belongs to and physically moving the data that changed ownership, all while the system keeps serving live traffic — effectively a full data migration under production load, not a configuration change.
What does "resharding" require operationally, and why is it considered risky and expensive?
It requires moving data between shards according to a new scheme, updating routing so queries find data at its new location, and doing all of this without losing writes or serving inconsistent reads mid-migration — usually via a careful dual-write or gradual-cutover process. It's risky because a mistake can corrupt or lose data across your entire dataset at once, not just one shard.
Why do cross-shard joins become expensive or impossible, and what's the usual workaround?
A join needs matching rows to be compared against each other, but if they live on different shards (different machines), the database can't just do it locally — it must fetch data from multiple shards and join it in the application or a coordinating layer instead. The usual workaround is denormalizing so related data needed together lives on the same shard, or avoiding true joins in favor of separate queries stitched together by the application.
What is a "scatter-gather" query, and why is it slower than a single-shard query?
A scatter-gather query is sent to every shard (or many shards) in parallel, and the results are collected and combined afterward. It's slower and more resource-intensive than a single-shard query because it multiplies the number of database round trips, and the overall response time is bounded by the slowest shard to respond, not the average.
Scenario: "how many total orders were placed today" across 8 shards — what does this actually require, and why is it more expensive than on an unsharded database?
It requires a scatter-gather query: asking all 8 shards for their own count and summing the results in the application (or a coordinator), rather than one query against one table. It's more expensive because it's now 8 separate queries plus a merge step, and it's only as fast as the slowest shard, whereas an unsharded database answers it with a single local count.
What is a distributed transaction across shards, and why is achieving atomicity harder there than within a single database?
It's a transaction whose writes span more than one shard (separate databases), needing all of them to either fully commit or fully roll back together. It's harder than a single-database transaction because each shard can independently succeed or fail, and there's no single, shared transaction log to atomically decide the outcome — the shards must explicitly coordinate to agree.
What is two-phase commit (2PC), briefly, and what's its main operational downside?
A coordinator asks every participating shard to "prepare" the transaction (confirming they're able to commit), and only once all of them agree does it tell everyone to actually commit. Its main downside is that if the coordinator crashes after some shards have prepared but before telling them to commit, those shards are left holding locks indefinitely, unable to safely commit or abort on their own.
Why do most sharded systems try to avoid needing cross-shard transactions altogether, rather than solve them with 2PC?
2PC adds real latency (extra coordination round trips on every transaction) and a real availability risk (a stuck coordinator can block shards indefinitely) — most systems find it far cheaper to design the shard key and data model so that transactions naturally stay within one shard, than to pay 2PC's cost on an ongoing basis.
What's the relationship between sharding and replication — are they mutually exclusive?
No — they're typically combined: each individual shard is usually also replicated on its own, so sharding solves the "too much data/write volume for one machine" problem while replication solves the "this one machine (per shard) shouldn't be a single point of failure" problem, independently of each other.
Scenario: a shard holding a very popular user's data on a social app becomes overloaded while other shards sit idle. What's this called, and how might you fix it without resharding everyone?
This is a hot shard (sometimes called a hotspot or "celebrity problem"). A targeted fix is isolating just that one heavily-loaded key onto its own dedicated shard (or splitting its data further, e.g. by content id instead of user id for that one user), rather than re-sharding the entire dataset to fix an issue caused by a small number of outlier keys.
Why can sharding by a monotonically increasing key, like an auto-incrementing id or a timestamp, create a hotspot?
New rows always have the newest (highest) value of that key, so under range-based sharding they all land in the same 'current' range on the same shard — every new write concentrates on whichever shard currently owns the highest range, no matter how many shards exist in total.
What's a composite shard key, and why might you use one instead of a single field?
A composite shard key combines more than one field (e.g. tenantId + userId) to compute shard placement. It's useful when a single field alone would distribute unevenly or wouldn't match your actual query patterns — combining fields can spread load more evenly while still letting related data (like everything for one tenant) stay groupable.
What is geo-sharding, and what benefit does it offer beyond raw scaling?
Geo-sharding assigns shards by geographic region, so a user's data lives on a shard physically close to them. Beyond distributing load, it reduces latency (data is closer to the user making requests) and can help satisfy data-residency requirements that require certain users' data to stay within a specific country or region.
Scenario: you need to move a single large tenant's data from an overloaded shard to a new, empty shard, without downtime. What does this actually require?
Typically: start dual-writing that tenant's new changes to both the old and new shard, backfill their existing historical data into the new shard, verify the copies match, then atomically flip routing for that tenant to the new shard and stop writing to the old one — all while the tenant keeps operating normally throughout.
What's the difference between a "local" secondary index and a "global" secondary index in a sharded system?
A local secondary index is built and maintained separately within each shard, covering only that shard's own data — cheap to maintain, but a query using it must still scatter-gather across every shard to get a complete answer. A global secondary index spans all shards' data in one structure, letting a query hit it directly without scatter-gather, but it now needs its own separate infrastructure and must stay in sync with writes happening across every shard.
Why is a global secondary index harder to maintain consistently than a local one?
Every write to any shard now also has to update this shared, cross-shard structure — which lives outside any single shard's transaction, so keeping the index and the underlying sharded data perfectly in sync (especially under concurrent writes and partial failures) is a genuinely harder distributed-consistency problem than a local index, which updates alongside its own shard's data in one place.
What happens to a unique constraint (like a unique email) when a table is sharded by a different key (like user id)?
The database can no longer enforce that uniqueness by itself, because two rows with the same email could land on two different shards that have no shared knowledge of each other's data. Enforcing it then requires either a separate global index/lookup service dedicated to checking uniqueness, or accepting weaker guarantees (like checking at write time and reconciling rare conflicts after the fact).
Why is capacity planning for shard count inherently an estimate, and what happens if you pick too few or too many?
It depends on projected data growth and traffic patterns that are genuinely hard to predict precisely. Too few shards and you'll hit the same scaling wall sharding was meant to solve, sooner than planned; too many shards and you pay unnecessary operational complexity and cross-shard query overhead for capacity you don't yet need — many systems deliberately over-provision logical shards on fewer physical machines up front, so future growth means moving logical shards rather than re-sharding the keyspace.
What's a common mistake teams make by sharding too early?
Taking on the real complexity of cross-shard queries, harder transactions, and operational overhead before the dataset or write volume has actually outgrown a single well-tuned, replicated database — often when simpler options (better indexing, caching, a bigger server, read replicas) would have bought significant additional headroom first.
Scenario: `country` is used as a shard key, but 80% of your users are in one country. What goes wrong, and what would you do differently?
One shard ends up carrying roughly 80% of the load while the rest sit comparatively idle — sharding in name only, since the actual bottleneck (that one shard) still exists. A better key would distribute load more evenly, e.g. hashing on user id, or further sub-sharding within that one dominant country rather than treating the whole country as a single unit.
Follow-up: what if you must keep `country` as the shard key for compliance reasons (data residency), despite the imbalance — what other levers do you have?
You can shard further within that one oversized country (e.g. by a hash of user id, still keeping all its shards within the required region/jurisdiction), give that region proportionally more physical shards than smaller regions get, or add caching/read replicas within that region to absorb load the shard key itself can't rebalance away.
Why does adding an index need to be done identically across every shard, and what happens if that gets out of sync?
Because each shard is a physically separate database, an index created on one isn't automatically created on the others — if it's applied inconsistently, queries can behave correctly (and fast) on some shards but slowly, or with different query plans, on others, making performance unpredictable depending purely on which shard happens to serve a given request.
How does sharding interact with a database's referential integrity (foreign keys) across shards?
A foreign key constraint is normally enforced by the database checking that the referenced row exists in the same database — but if the referencing and referenced rows can end up on different shards, there's no single database that can enforce that check anymore, so referential integrity across shards typically has to be maintained by application logic instead of the database itself.