Databases at Scale
When one database isn't enough: read replicas, sharding and partitioning, SQL vs. NoSQL, and the CAP theorem's forced choice between consistency and availability under a partition.
On this page
Caching absorbs reads, but eventually the database itself becomes the bottleneck - too many writes, too much data for one machine, or a need to survive failures. Scaling a database is where system design gets genuinely hard, because it forces you into fundamental trade-offs you cannot escape: consistency versus availability, normalization versus speed, one big machine versus many coordinated ones. This lesson maps the tools and the trade-off that governs them all.
Scale up, then scale out
Two directions to add capacity:
- Vertical scaling (scale up) - a bigger machine: more CPU, RAM, faster disks. Simple, no code changes, and surprisingly far-reaching - but there's a ceiling (the biggest machine money can buy) and a single point of failure.
- Horizontal scaling (scale out) - more machines working together. Effectively unlimited, and fault-tolerant, but it introduces the coordination complexity that makes distributed data hard.
You scale up until it's uneconomical or impossible, then scale out. The techniques below are all forms of scaling out.
Replication: copies for reads and resilience
Replication keeps copies of the data on multiple machines. The common pattern is leader-follower (primary-replica): writes go to the leader, which streams changes to read-only followers. This buys two things:
- Read scaling - route the flood of reads (recall: most systems are read-heavy) to many followers, while the leader handles writes.
- High availability - if the leader dies, a follower is promoted. Copies survive machine failure.
The catch is replication lag: a follower may be milliseconds-to-seconds behind the leader, so a read from a replica can return slightly stale data (you just wrote something, then read it from a lagging follower and don't see it). That staleness is the price of read-scaling - usually fine, occasionally not (read-your-own-writes needs care).
Sharding: splitting the data itself
Replication copies the whole dataset; when the data is too big for one machine (or writes exceed one leader's capacity), you must shard (partition) - split the data across machines, each holding a subset:
Shard by a key (e.g. user_id):
users A-H → shard 1
users I-P → shard 2 each shard is an independent DB holding part of the data
users Q-Z → shard 3Sharding scales writes and storage (each shard handles its slice), but it's the hardest tool here:
- Choosing the shard key is critical. A bad key creates hot shards (one shard gets most traffic - e.g. sharding by country when one country dominates). A good key spreads load evenly (often a hash of the id).
- Cross-shard queries are painful. A query spanning shards (join, aggregate) must hit many shards and combine results - slow and complex. You design shards so common queries stay within one shard.
- Rebalancing is hard. Adding a shard means moving data; consistent hashing minimizes the movement.
The CAP theorem: the unavoidable choice
Underneath distributed data sits an inescapable law. The CAP theorem states that in the presence of a network partition (nodes can't all talk to each other), a distributed system must choose between:
- Consistency (C) - every read sees the latest write (all nodes agree), or
- Availability (A) - every request gets a response, even if some data might be stale.
You cannot have both during a partition - you must pick. (When there's no partition, you can have both.)
- CP systems (choose consistency) - refuse or block requests rather than serve possibly-stale data. Good when correctness is paramount (a bank balance). Examples: many traditional SQL setups, HBase, ZooKeeper.
- AP systems (choose availability) - keep serving, accepting possible staleness, and reconcile later (eventual consistency). Good when uptime matters more than perfect freshness (a social feed, a shopping cart). Examples: Cassandra, DynamoDB.
This is why the choice isn't "SQL vs NoSQL" in the abstract - it's what does this data need? Money needs consistency (CP); a like-count can be eventually consistent (AP).
SQL vs NoSQL, through this lens
The SQL/NoSQL choice is really about which trade-offs a system makes:
- Relational (SQL) - strong consistency, ACID transactions, flexible queries and joins, a rigid schema. Great for complex, related, correctness-critical data (the ledger!). Scales up easily; scales out (sharding) with effort.
- NoSQL - a family (key-value, document, wide-column, graph) that typically trades joins and strict schema for easy horizontal scaling and often AP/eventual consistency. Great for huge scale, simple access patterns, or flexible schemas.
The senior answer is rarely "always NoSQL" - it's matching the store to the data's consistency, query, and scale needs. Many systems use both: SQL for the transactional core, NoSQL for the high-scale, simple-access parts.
Don't shard until you must - it's a one-way door of complexity
Sharding solves a real problem (data or writes exceeding one machine) but imposes lasting costs: cross-shard queries, distributed transactions, rebalancing, and a shard key you're stuck with. Exhaust the cheaper options first - vertical scaling, read replicas, caching, and archiving old data. Many systems that 'need sharding' really need a bigger box and a cache. When you do shard, choose the key with extreme care; changing it later means re-sharding the whole dataset.
One database is a single library building. First you scale up - taller shelves, a bigger building - until you hit the biggest building possible. Replication is opening identical branch libraries with copies of every book: readers (reads) spread across branches, and if one burns down the collection survives - but a book just donated to the main branch takes time to copy to the others (replication lag). Sharding is splitting the one collection across buildings by subject - science here, history there - so no single building holds everything; brilliant until you need a book that spans subjects (cross-shard query) or one subject becomes wildly popular (hot shard). And CAP is the moment the phone lines between branches go down (a partition): do you refuse to lend a book because you can't confirm it wasn't checked out elsewhere (consistency), or lend it anyway and sort out duplicates later (availability)? You can't do both while the lines are down - you must choose.
A payments system stores two kinds of data: (1) account balances and transactions - must be exactly correct, a lost or double-counted transaction is catastrophic; (2) a feed of 'recent activity' notifications shown in the UI - high volume, and it's fine if it's a few seconds stale or occasionally shows a duplicate. For each, choose SQL vs NoSQL and CP vs AP, and name one scaling technique you'd apply.
What does the CAP theorem force you to choose between during a network partition?
Key takeaways
- Scale up (bigger machine) until it's uneconomical, then scale out (more machines) - replication and sharding are forms of scaling out.
- Replication (leader-follower) copies the whole dataset for read-scaling and high availability, at the cost of replication lag (stale reads from followers).
- Sharding splits the data across machines to scale writes and storage; the shard key is critical (avoid hot shards) and cross-shard queries are painful.
- The CAP theorem: during a network partition you must choose Consistency (CP) or Availability (AP) - money needs CP, a social feed can be AP.
- SQL vs NoSQL is really a trade-off choice (ACID/consistency/joins vs easy horizontal scale/eventual consistency); match the store to the data, and often use both.