Skip to content
Tech Interview Prep home
Technical interview guide

Database Scaling (Sharding & Replication)

Splitting data across machines (sharding) and copying it across machines (replication) — solving two different scaling problems.

Read
24 min
Practice MCQs
25
Interview QA
25
Edition
v3
Editorial status
Reviewed

Scope: Vendor-neutral distributed-systems principles; current AWS architecture guidance.

Overview

Curated: · Written: · Reviewed:

Scale the thing that is actually saturated

Database scaling is a sequence of increasingly expensive and irreversible steps, and the discipline is refusing to skip ahead. Almost every system that "needs sharding" needs an index, and almost every system that needs a bigger machine needs its top three queries fixed. The ladder below is ordered by cost and regret, and each rung should be exhausted before the next.

Before scaling anything: find the saturated resource

Which is it — CPU, disk I/O, memory, connections, or lock contention? Without that, every remedy is a guess, and the common ones point in different directions. An I/O-bound system gains nothing from more cores. A connection-exhausted system gains nothing from more memory.

Then look at the workload. Aggregate query statistics almost always show that a handful of queries dominate total time, and total time is the right ranking — a fast query called a million times consumes more capacity than a slow one called twice. Fixing the top two or three routinely reclaims more headroom than a hardware tier, and costs nothing ongoing.

Check maintenance too. Failing vacuum causes bloat, which inflates every scan and destroys cache efficiency — and it looks exactly like insufficient memory from the outside. A single long-running transaction can hold cleanup back database-wide. These are free to fix and are frequently the whole problem.

Rung 1: query and index optimisation

The cheapest and most reversible lever. A missing index on a foreign key, a predicate wrapping a column in a function so the index cannot be used, SELECT * preventing an index-only scan, deep OFFSET pagination, an N+1 pattern issuing one query per row — each of these can remove substantial work without an architectural redesign, though an added index still carries write and storage cost.

Indexes are not free: each is maintained on every write. A well-chosen index often beats a hardware upgrade for the query it serves because it reduces work, but its benefit and write cost must be measured.

Rung 2: connection pooling

A database has a hard limit on useful concurrent connections, and exceeding it makes things worse rather than slower — context switching and per-connection memory dominate. A pooler multiplexes many application connections onto few database ones, and adding one often produces a step change that no amount of hardware would.

This is also where a great deal of apparent database pressure turns out to originate: long transactions holding connections, so the pool exhausts and requests queue waiting for a connection while the database itself is barely busy.

Rung 3: caching

A cache hit removes that read from the database rather than merely serving the same database query faster; misses and invalidation work remain. It suits data read far more often than it changes, and it is the highest-leverage step for a read-heavy workload.

The cost is invalidation, and there are only three honest strategies: a TTL, which bounds staleness and cannot get permanently out of step; explicit invalidation on write, which is fresher but fails silently when a write path forgets; and write-through, which centralises update logic but still needs failure handling to avoid database/cache divergence and may couple writes to cache availability. TTL plus explicit invalidation is a good default, because the TTL bounds the damage when an invalidation is missed — and one eventually will be.

Rung 4: vertical scaling

More CPU, memory or faster storage. It is genuinely the right answer more often than its reputation suggests: it requires no application change, no consistency compromise, and modern machines are very large. If the working set no longer fits in memory and the queries are already efficient, more memory is exactly right.

Its limits are a ceiling, cost that grows non-linearly at the top end, and usually a restart to resize.

Rung 5: read replicas

Replicas serve read traffic from copies, which suits read-heavy workloads — most workloads. Analytics and reporting move off the primary entirely, removing both CPU and cache-eviction pressure.

The cost is replication lag, and it breaks assumptions rather than performance. Read-your-writes is the one that produces visible bugs: a user updates something, the next read hits a lagging replica, and they see the old value and conclude the save failed. Remedies are routing reads to the primary for a window after a write, waiting for a known replication position, or — often simplest and most overlooked — returning the written value from the write response instead of re-fetching it.

Replicas do not help write throughput at all. Every write still goes to the primary and is replicated to every replica, so adding replicas adds write work.

Rung 6: partitioning within one database

Splitting a large table into partitions by a key — usually time. Queries filtering on that key prune to relevant partitions, indexes are per-partition and therefore smaller, and retention becomes cheap: detaching a monthly partition is near-instant, where deleting a month of rows is expensive and leaves bloat.

The constraints must be accepted deliberately: any unique constraint must include the partition key; queries not filtering on it touch every partition; and future partitions must be created automatically, since a missing one causes insert failures at a period boundary.

Rung 7: sharding

Splitting data across independent databases by a shard key. Within this ladder it is the step that distributes write ownership horizontally, and it is by far the most expensive and least reversible; distributed-SQL or multi-writer architectures are separate designs with their own coordination costs.

What you give up is substantial: cross-shard joins and transactions, so anything spanning shards moves into application code; unique constraints across shards; and simple operations — every migration, backup and query now runs N times. Rebalancing when a shard grows too large is a genuinely hard operation.

Shard key choice is the decision that matters and cannot easily be changed. It must distribute evenly, be present in nearly every query so requests route to one shard, and keep related data co-located so common operations stay single-shard. Sharding by customer usually satisfies all three; sharding by time concentrates all current writes on one shard.

Before sharding, consider functional partitioning — moving distinct workloads to separate databases by domain rather than splitting one dataset. It gives much of the relief with far less complexity, and it is frequently sufficient.

What to hold on to

Two facts worth keeping in view. Most systems never exceed what a well-indexed relational database on one large machine handles, and that number is much higher than people assume. And every step down this ladder trades away something — consistency, operational simplicity, or query flexibility — permanently, in exchange for capacity you may not need. Measure before each step, and be willing to stop.

Worked example: four replicas do not raise write tps

Primary saturates at 1,200 write tps. Checkout filter is WHERE tenant_id = $1 with no index. p99 read is 22 s.

changewrite tpscheckout p99
4 read replicas1,20022 s (writes still primary; plan still seq scan)
btree on tenant_id1,2009 ms
shard by created_atcurrent month still one hot shardwrites unchanged

Replicas scale reads. An index, then functional split, then a shard key that is in every write — in that order.