Scaling Beyond One Machine
A single database node eventually hits a wall: CPU, RAM, IOPS, or connection limits. Scaling has two directions. Vertical scaling (a bigger box) is simple but bounded and expensive at the top. Horizontal scaling (more boxes) is unbounded but forces you to confront replication, partitioning, and distributed consistency. This page is about the horizontal world.
The order you reach for tools
Indexing → caching → read replicas → vertical scale → partitioning/sharding → distributed SQL. Sharding is the last resort because it complicates joins, transactions, and operations.
SQL vs NoSQL, ACID vs BASE
SQL / relational systems offer a fixed schema, joins, and ACID transactions. NoSQL trades some of these for horizontal scale and flexible schemas, often under BASE semantics.
| ACID | Meaning |
|---|---|
| Atomicity | All-or-nothing; a transaction fully commits or fully rolls back |
| Consistency | Constraints/invariants hold before and after |
| Isolation | Concurrent transactions don't corrupt each other |
| Durability | Committed data survives crashes (WAL/fsync) |
BASE = Basically Available, Soft state, Eventually consistent. The system stays available and converges over time rather than guaranteeing an immediately consistent view. It is a design stance, not a weaker ACID.
NoSQL families
| Type | Model | Examples |
|---|---|---|
| Key-value | Opaque value by key | Redis, DynamoDB, Riak |
| Document | JSON-like nested docs | MongoDB, Couchbase |
| Wide-column | Sparse rows, column families | Cassandra, ScyllaDB, HBase |
| Graph | Nodes + edges, traversals | Neo4j, Neptune |
Normalization vs Denormalization
Normalization removes redundancy so each fact lives in one place — great for write integrity, painful for read-heavy joins at scale. Denormalization duplicates data to serve reads without joins — fast reads, but writes must update every copy and can drift. At scale you often denormalize deliberately and keep copies in sync with events/CDC.
Indexing
| Index | Strength | Used by |
|---|---|---|
| B-tree | Range + equality, sorted scans | PostgreSQL, MySQL (InnoDB) |
| Hash | Exact-match O(1); no ranges | Memory tables, hash indexes |
| LSM-tree | Write-optimized; sequential flushes + compaction | Cassandra, RocksDB, ScyllaDB |
| Inverted | Full-text / term lookups | Elasticsearch, Lucene |
B-trees update in place (read-optimized, good for mixed workloads). LSM-trees buffer writes in memory and flush sorted runs to disk, so writes are cheap and sequential but reads may check several files (mitigated with Bloom filters), and background compaction costs I/O.
Replication
Leader-Follower (single-leader) Leaderless (quorum)
writes W nodes ack write
| R nodes serve read
[Leader] --replicate--> [Follower] W + R > N => overlap
| \-> [Follower] (N replicas total)
reads/writes reads only
| Topology | Writes | Trade-off |
|---|---|---|
| Single-leader | One leader accepts writes | Simple; leader is a write bottleneck & SPOF (failover needed) |
| Multi-leader | Several leaders (per region) | Low write latency per region; write conflicts to resolve |
| Leaderless | Any replica; quorum | Highly available; needs read-repair / anti-entropy |
Synchronous replication waits for a follower to confirm (no data loss on leader failure, higher latency). Asynchronous replication returns immediately (fast, but a leader crash can lose the last un-replicated writes and readers see replication lag). Quorum systems use the rule W + R > N so read and write sets overlap, guaranteeing a read sees the latest write.
Read-your-writes
Async read replicas can serve stale data. If a user updates their profile then re-reads from a lagging replica, they see the old value. Fixes: route the user's reads to the leader for a short window, or track a write timestamp/version and only read from replicas caught up to it.
Partitioning / Sharding
Partitioning splits one logical dataset across nodes so each holds a subset. The partition key choice is the single most important scaling decision — a bad key creates hot shards.
| Strategy | How | Downside |
|---|---|---|
| Range | By key ranges (A–M, N–Z) | Hot spots on sequential keys |
| Hash | By hash(key) | Even spread; loses range queries |
| Directory | Lookup table maps key→shard | Flexible; directory is a SPOF/bottleneck |
Rebalancing moves partitions when you add/remove nodes. Avoid hash % N (adding a node reshuffles everything). Consistent hashing or a fixed number of partitions (far more than nodes, reassigned as a unit) keeps rebalancing cheap.
Consistent hashing ring [0, 2^32):
key -> next node clockwise; add/remove moves only ~1/N of keys.
Virtual nodes give each physical node many ring points => even load.
(Cassandra, DynamoDB, Riak, Vitess use variants of this.)
CAP in Databases
Under a network Partition you must choose Consistency (reject requests to avoid divergence) or Availability (serve possibly stale data). It is a choice made during partitions; the healthy path can be both fast and consistent. PACELC extends it: else, absent partitions, you still trade Latency vs Consistency.
| Lean | Systems |
|---|---|
| CP (consistency) | HBase, MongoDB (majority), Spanner, etcd, ZooKeeper |
| AP (availability) | Cassandra, DynamoDB, Riak (tunable) |
Isolation Levels & Anomalies
| Level | Dirty read | Non-repeatable | Phantom |
|---|---|---|---|
| Read Uncommitted | Possible | Possible | Possible |
| Read Committed | No | Possible | Possible |
| Repeatable Read | No | No | Possible* |
| Serializable | No | No | No |
Dirty read: reading another transaction's uncommitted write. Non-repeatable read: the same row returns different values within one transaction. Phantom: a range query returns different rows because another transaction inserted/deleted. *MySQL InnoDB's Repeatable Read blocks phantoms via next-key locks; the SQL standard permits them.
Distributed Transactions: 2PC vs Saga
| Two-Phase Commit | Saga | |
|---|---|---|
| Model | Coordinator: prepare then commit | Sequence of local txns + compensations |
| Atomicity | Strong (all or none) | Eventual; roll back via compensations |
| Blocking | Locks held; blocks if coordinator dies | Non-blocking, per-step commit |
| Fit | Few nodes, short txns | Microservices, long-running flows |
Operational Building Blocks
Connection pooling (PgBouncer, HikariCP) reuses a bounded set of DB connections; databases handle far fewer connections than app instances can open. Read replicas offload read traffic (watch lag). Change Data Capture (CDC) streams the commit log (Debezium reading the WAL/binlog) to keep caches, search indexes, and downstream services in sync without dual-write bugs.
CDC pipeline:
Postgres WAL / MySQL binlog --> Debezium --> Kafka --> [search index]
\--> [cache]
\--> [data warehouse]
NewSQL & OLTP vs OLAP
NewSQL aims for the holy grail: horizontal scale and SQL + ACID. Google Spanner uses TrueTime (GPS/atomic clocks) for external consistency; CockroachDB and YugabyteDB layer distributed SQL over a Raft-replicated KV store; Vitess shards MySQL transparently (powers YouTube, Slack).
OLTP (transactional) systems handle many small, low-latency reads/writes (row-oriented, e.g. Postgres). OLAP (analytical) systems scan huge volumes for aggregates (column-oriented, e.g. Snowflake, ClickHouse, BigQuery). Don't run heavy analytics on your OLTP primary — ETL/CDC the data into a warehouse.
Capacity sketch
10 TB of data, ~1 TB usable per shard for headroom → ~10 shards. At 200k writes/s with each shard sustaining ~30k writes/s, you need ~7 shards for throughput — take the max (10). Add replicas per shard for HA and read scaling, then size the connection pool to the DB's real limit, not the app fleet size.