Scaling a database, in the right order

JR

Jai Rao

August 24, 202610 min read

Indexes, pooling and a bigger machine before replicas, partitions and sharding. Why hash-modulo sharding moves 75 percent of your data when you add a node.


Sharding is the last resort, not the goal. Most teams who believe they need it need an index, and a surprising number of the rest need a connection pool. This is worth saying up front because the topic attracts aspiration — sharding sounds like what serious systems do — and the aspiration leads people to take on years of operational cost to solve a problem that a five-minute change would have fixed.

So here is the ladder, in the order you should climb it, with the honest price of each rung. Do not skip to the bottom.

Find out what is actually slow

Before any of it: measure. "The database is slow" is not a diagnosis, and the fix for a missing index has nothing in common with the fix for connection exhaustion.

Read the query plan. Every database will tell you how it intends to execute a query, and the thing to look for is a sequential scan where you expected a lookup, plus the ratio of rows examined to rows returned. Examining two million rows to return ten is the signature of a missing index, and it is visible before you optimise anything.

Text
-- Illustrative plan, not captured from a specific system.EXPLAIN SELECT * FROM orders WHERE customer_id = 4711;Seq Scan on orders  (rows removed by filter: 1999990)  Filter: (customer_id = 4711)-- examined 2,000,000 rows to return 10CREATE INDEX idx_orders_customer ON orders (customer_id);Index Scan using idx_orders_customer on orders  (rows: 10)  Index Cond: (customer_id = 4711)-- examined 10 rows to return 10

The reason that difference is so large is worth understanding rather than accepting. An index is a B-tree, and a B-tree page holds many keys — on the order of a hundred. So each level you descend narrows the search by a factor of a hundred, which means the depth grows logarithmically while the table grows linearly:

Text
          1,000 rows   full scan         1,000   B-tree ~  1.5 levels      1,000,000 rows   full scan     1,000,000   B-tree ~  3.0 levels    100,000,000 rows   full scan   100,000,000   B-tree ~  4.0 levels

A hundred million rows is about four page reads. That is the whole argument for indexing, and it also explains why the gap widens as you grow: the scan cost multiplies by a hundred while the lookup cost goes up by one.

Two traps, both common. Indexes are not free — every one must be updated on every write, so a table with twelve indexes has slow writes by construction, and unused indexes are pure cost. An index the planner cannot use is worse than none, because you believe you are covered. Wrapping the column in a function (WHERE lower(email) = ...) or leading a pattern with a wildcard (LIKE '%smith') both defeat a plain B-tree index. Check the plan rather than the intention.

The cheap rungs everyone skips

Connection pooling. A database serves connections with a process or thread and a memory allocation each; it wants dozens, not thousands. Point two hundred application instances at it with ten connections apiece and you have asked for two thousand — and the failure is not graceful degradation, it is refused connections and an outage caused entirely by your own client configuration. A pooler in front, multiplexing many client connections onto few database ones, is a configuration change that resolves an entire class of incident.

The N+1 query. Fetch fifty orders, then loop and fetch each customer: fifty-one round trips where two would do. Individually every query is fast and correctly indexed, which is what makes this hard to spot from database metrics — nothing looks slow, there is just an enormous amount of it. This is the most common cause of a slow page in an application backed by an ORM.

Buy a bigger machine. Unglamorous, effective, and worth arguing for out loud. Doubling memory and CPU is an afternoon of work with a predictable outcome and no new failure modes. Sharding is months of work that permanently complicates every future feature. There is a real point where vertical scaling stops being available, but plenty of teams take on distributed-systems complexity while still an order of magnitude below the largest single instance they could rent.

Read replicas, and the comment that vanishes

Most workloads read far more than they write, so copying the data to replicas and sending reads there is the natural next step. The primary takes writes and streams its changes; replicas apply them and serve reads.

The complication is that streaming takes time. Replication lag is usually milliseconds and occasionally much worse — during a write burst, a long transaction, or a network problem — and it produces a bug that is nearly impossible to reproduce on a developer machine where lag is zero.

A user posts a comment. The write goes to the primary. The page reloads, the read goes to a replica that has not caught up, and the comment is not there. From the user's point of view their submission was lost, so they post it again. Three approaches, in ascending order of precision:

  • Read from the primary after writing. For a short window after a user's write, send that user's reads to the primary. Crude, simple, and effective for the common case.
  • Pin the session. Route a user consistently to one replica so they at least see a monotonic view — they may see slightly old data, but never data that goes backwards.
  • Carry a version token. Record the write position, pass it with the read, and wait for the replica to reach it (or fall back to the primary). Precise, and the most work.

Whichever you choose, decide it deliberately per read path. "Send all reads to replicas" is a configuration that works in testing and generates support tickets in production. And note that replicas do not help write throughput at all — if writes are your bottleneck, this rung is not your rung.

Partitioning is not sharding

These get conflated, and the difference is the whole point. Partitioning splits one table into pieces inside one database. Sharding splits data across independent machines. Partitioning is a local storage optimisation and costs you almost nothing; sharding is a distributed system and changes what your application can do.

Partitioning shines on time-series data. Partition by month, and a query for last week touches one partition instead of the whole table. Better still, expiring old data becomes dropping a partition — an instant metadata operation — rather than a DELETE of a hundred million rows that generates enormous write load and leaves the table needing a cleanup.

If your access pattern has a natural range dimension, especially time, try this before considering sharding. It is frequently the last rung you need.

Sharding, and the decision you cannot reverse

Now the real thing. Data is split across machines, each holding a subset, and the application must know where to look.

The shard key is the decision that matters, because changing it later means moving all of your data while serving traffic. Get it wrong and you get hotspots — the whole point is spreading load, and a bad key concentrates it:

  • Sharding by tenant works until one customer is twenty times bigger than the rest, and their shard is permanently on fire while others idle. Very common in business software, where customer sizes follow a power law.
  • Sharding by timestamp sends every new write to whichever shard owns the current period, so one shard takes all the write traffic and the others are read-only archives.
  • Sharding by a monotonic id has the same defect for the same reason.

You want a key with high cardinality, even distribution, and — critically — one that appears in most of your queries. A perfectly distributed key that your queries do not filter on means every query hits every shard, which is worse than not sharding.

Be clear-eyed about what you give up:

  • Joins across shards stop being a database operation and become application code that fetches from several places and combines results.
  • Transactions across shards require distributed protocols you almost certainly do not want to operate. In practice you redesign so transactions stay within one shard.
  • Global uniqueness is no longer enforceable by a constraint. Unique email addresses need a separate authority or a scheme that partitions the namespace.
  • Aggregate queries — count, sum, group by — must run everywhere and be merged, so "how many users signed up today" changes from a one-line query into a small pipeline.

Consistent hashing, and why modulo is a trap

The obvious way to map a key to a shard is to hash it and take the remainder over the shard count. It distributes evenly and it is a trap, because the shard count is baked into every assignment. Add one shard and almost every key belongs somewhere else.

Measured over twenty thousand keys, going from three shards to four:

Text
naive hash % node_count: 3 -> 4  keys moved: 14,892 of 20,000 = 74.5%

Three quarters of your data has to move, while the system stays online. That is not a maintenance window, it is a migration project — and it recurs every time you grow.

Consistent hashing fixes this by hashing nodes onto a ring alongside keys, so a key belongs to the next node clockwise. Adding a node only captures the keys between it and its predecessor; everything else stays put. Several virtual points per physical node keep the distribution even:

Text
import hashlib, bisectdef build_ring(nodes, vnodes=150):    """Each node appears many times so ownership is evenly spread."""    ring = []    for n in nodes:        for v in range(vnodes):            h = int(hashlib.md5(f"{n}#{v}".encode()).hexdigest()[:8], 16)            ring.append((h, n))    ring.sort()    return ringdef owner(ring, key):    h = int(hashlib.md5(key.encode()).hexdigest()[:8], 16)    i = bisect.bisect(ring, (h,))    return ring[i % len(ring)][1]      # next node clockwise

The same growth, the same keys:

Text
consistent hashing: 3 nodes -> 4 nodes  keys moved: 4,873 of 20,000 = 24.4%   (ideal 1/4 = 25.0%)  distribution after: {'n1': 4702, 'n2': 5058, 'n3': 5367, 'n4': 4873}

Twenty-four percent moved, against a theoretical minimum of twenty-five — essentially optimal, and only the data that genuinely must relocate. The distribution is even to within a few percent, which is what the virtual nodes buy: with one point per node the ring would be lumpy and some shard would own twice its share by accident.

A checklist keyed to symptoms

Work down this list and stop at the first thing that matches.

  • A few specific queries are slow, others are fine. Missing or unusable index. Read the plan.
  • Everything is slow, including trivial queries. Connection exhaustion or CPU saturation. Check the pool before touching schema.
  • One page is slow but every query in it is fast. N+1. Count the queries per request.
  • Reads are heavy, writes are modest. Replicas — and decide the read-your-writes strategy at the same time.
  • Queries filter on a date range, old data is rarely touched. Partition by time.
  • Write throughput exceeds one machine, and vertical scaling is genuinely exhausted. Now consider sharding.

What would justify that last step: writes saturating the largest instance available to you, a working set that cannot fit in any single machine's memory, or a regulatory requirement to hold data in separate places. Those are real reasons. "We expect to grow" is not one, because you can shard later with more information and fewer wrong guesses — and the shard key you would choose today, before seeing real access patterns, is likely to be the one you regret.

Whatever rung you are on, keep watching the same handful of numbers: slow query log, rows examined against rows returned, connection pool utilisation and wait time, replication lag, and cache hit ratio in the buffer pool. Those five will tell you which rung you are actually on, which is usually a lower one than it feels like.