Databases

Database sharding: interview questions and how to answer them

Sharding splits data across independent databases to scale writes and storage, at the cost of cross-shard joins, transactions and operational simplicity.

Written and reviewed by Sahil Srivastav

DatabasesScalingLast resort

What it actually is

Sharding distributes rows across multiple independent database servers according to a shard key. Each shard holds a disjoint subset and knows nothing about the others, so writes and storage scale roughly linearly with shard count — which is the one thing no amount of indexing, caching or read replicas will give you.

The reason it is a last resort is that it removes guarantees the rest of your system was built on. A join across shards is no longer a join; it is application code fetching from several databases and combining results. A transaction across shards does not exist; you need a saga with compensations, or a design where it never arises. Unique constraints are per-shard. Even ORDER BY ... LIMIT across shards means over-fetching from each and merging.

The decision that determines everything afterwards is the shard key, because it decides which queries stay cheap. Queries carrying the key hit one shard. Queries without it must fan out to all of them, and the slowest shard sets the latency. The key is also expensive to change later, which is why interviews concentrate on how you would choose it rather than on the mechanics of routing.

Why it matters in production

Because it is the correct answer to exactly one problem — a single primary that cannot absorb the write volume or hold the data — and the wrong answer to most others. Read load is solved by replicas and caching. Query latency is usually solved by indexing and access-pattern fixes. Operational pain on huge tables is solved by partitioning. Reaching for sharding before exhausting those multiplies operational cost for no benefit, which is why "why not shard?" is as revealing an interview question as "how would you shard?".

It matters in interviews specifically because candidates who have not operated a sharded system tend to describe only the routing. The follow-ups are always about what broke: the query that needed a cross-shard join, the tenant that outgrew its shard, the migration that had to move live data without downtime.

How it works

Shard key choice decides your query costs

A key present in nearly every query — tenant, user, account — keeps most reads single-shard. A key absent from common queries forces fan-out, where every request hits every shard and tail latency dominates. The test is to list your top ten queries and count how many carry the candidate key.

Range, hash and directory, and what each trades

Range keys keep ordered scans efficient but create hotspots when traffic is skewed toward recent data. Hashing distributes evenly and destroys range queries. A directory — an explicit lookup table mapping key to shard — gives full control and the ability to move individual tenants, at the cost of a lookup on the critical path and a component that must not fail.

Resharding is the hard part

With naive modulo hashing, changing shard count remaps nearly every key. Consistent hashing limits movement to roughly 1/n of the data. The usual practical answer is to over-provision logical shards up front — say 1024 — and map many logical shards to each physical server, so growth means moving whole logical shards rather than rehashing anything.

Cross-shard operations have no good general answer

Joins become application-side fan-out and merge. Transactions become sagas with compensating actions. Global aggregates become per-shard partial aggregates combined afterwards. Global uniqueness needs either a key containing the shard key, or a separate allocation service. Each of these is a permanent tax on every feature that needs it.

Skew is the failure that actually happens

Shards are rarely equally loaded, because tenants are not equally sized. One customer ten times larger than the rest makes their shard the bottleneck regardless of how many shards exist. Systems that survive this have a way to relocate a single tenant — which is the strongest argument for directory-based mapping.

Implementing it

Exhaust the alternatives first and be able to say you did: read replicas, caching, indexing, archiving, partitioning, and a larger instance. Sharding is justified by a write or storage ceiling, not by general slowness.

Choose the key by auditing real queries, not by picking the most "natural" entity. Count how many of your top queries carry it; if it is not most of them, the key is wrong.

Plan resharding before you need it. Many logical shards mapped onto few physical servers is the cheapest form of future-proofing, and it costs almost nothing up front.

Keep a per-shard load and size metric from day one. Skew is the normal state, and you want to see it developing rather than discover it when one shard saturates.

// Naive modulo: adding a shard remaps nearly every key. Avoid.
const shard = hash(tenantId) % SHARD_COUNT;

// Logical shards fixed forever; physical mapping is data you can change.
// Moving a tenant = moving one logical shard, no rehashing.
const LOGICAL_SHARDS = 1024;
const logical = hash(tenantId) % LOGICAL_SHARDS;
const physical = shardMap.lookup(logical);   // directory, cached

// Single-shard: carries the shard key. Cheap.
//   SELECT * FROM orders WHERE tenant_id = $1 AND created_at > $2
//
// Fan-out: no shard key. Hits every shard; the slowest one sets latency,
// and LIMIT means over-fetching from each and merging in the application.
//   SELECT * FROM orders WHERE status = 'PENDING' ORDER BY created_at LIMIT 20

Interview questions and how to answer them

When should you shard?

When a single primary cannot absorb the write volume or hold the data, and you have already exhausted replicas, caching, indexing, archiving and partitioning. It is the only one of those that scales writes, and it is also the only one that removes joins and transactions — so the bar for reaching it should be high and evidence-based.

How do you choose a shard key?

By auditing the actual query mix. A key present in nearly every query keeps requests single-shard; one that is absent forces fan-out where tail latency dominates. Also weigh distribution — a key that concentrates traffic on one shard reproduces the original bottleneck — and whether you can relocate an individual value when one grows disproportionately.

What happens to transactions?

Within a shard they are normal. Across shards they do not exist: there is no distributed transaction available, so you use a saga with compensating actions and accept intermediate states, or you design so that anything needing atomicity lives on one shard. A candidate proposing two-phase commit here is describing something that blocks on coordinator failure and is rarely acceptable.

How do you add a shard without downtime?

Not by rehashing everything, which is why the mapping design matters more than the hash. With many logical shards mapped to few physical servers, you move whole logical shards: start dual-writing, backfill, verify, cut reads over, stop dual-writing. With consistent hashing, roughly 1/n of keys move. With naive modulo, almost everything moves, which is effectively a migration project.

One tenant is far larger than the others. What do you do?

Relocate them, which requires a mapping that supports moving a single key — a directory is best suited to this. If they are too large for any single shard, they need sub-sharding on a secondary key or a dedicated deployment. Hash-only mappings make this very difficult, which is the practical argument against pure hashing at scale.

Answers that lose the round

  • Sharding to fix read latency, which replicas and indexes solve far more cheaply
  • Choosing a shard key absent from most queries, so everything fans out
  • Naive modulo hashing, making the first reshard a full data migration
  • Assuming transactions still work across shards
  • No plan for a tenant that outgrows its shard
  • Treating sharding and partitioning as interchangeable
  • No per-shard metrics, so skew is discovered as an outage

Practise in a real repository

Explaining a concept and enforcing it in code are different skills, and machine coding rounds test the second. Gronex ships broken backend repositories whose test suites assert the invariant rather than the happy path.

FAQ

Is sharding the same as partitioning?

No. Partitioning splits a table within one database, keeping transactions, joins and constraints intact. Sharding splits data across separate databases and gives those up in exchange for write and storage scaling. Partitioning is dramatically cheaper and should be exhausted first.

Do managed distributed databases remove the problem?

They move it. Systems that distribute automatically still have a partition key that determines whether a query is local or fans out, and skew still hurts. The vocabulary changes; the shard-key decision does not go away.

How do you generate unique ids across shards?

Not with a per-shard sequence. Options are a key containing the shard identifier, a time-ordered scheme like Snowflake that embeds a node id, or UUIDs — noting that random UUIDs as a clustered primary key hurt insert locality, which is why time-ordered variants are preferred.

Can you un-shard later?

Technically yes, by consolidating onto a larger instance, and it is occasionally the right call once hardware improves or data is archived. It is a migration of the same magnitude as sharding was, so it is rarely done — which is another reason to delay the original decision until the evidence is unambiguous.

Related

More backend concepts