Learn/system design/Databases at Scale
Advanced~25 min read

Databases at Scale

Replication, partitioning, isolation levels, distributed transactions, and NewSQL for databases beyond one machine.

ShardingReplicationCAPNewSQL

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.

ACIDMeaning
AtomicityAll-or-nothing; a transaction fully commits or fully rolls back
ConsistencyConstraints/invariants hold before and after
IsolationConcurrent transactions don't corrupt each other
DurabilityCommitted 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

TypeModelExamples
Key-valueOpaque value by keyRedis, DynamoDB, Riak
DocumentJSON-like nested docsMongoDB, Couchbase
Wide-columnSparse rows, column familiesCassandra, ScyllaDB, HBase
GraphNodes + edges, traversalsNeo4j, 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

IndexStrengthUsed by
B-treeRange + equality, sorted scansPostgreSQL, MySQL (InnoDB)
HashExact-match O(1); no rangesMemory tables, hash indexes
LSM-treeWrite-optimized; sequential flushes + compactionCassandra, RocksDB, ScyllaDB
InvertedFull-text / term lookupsElasticsearch, 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

text
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
TopologyWritesTrade-off
Single-leaderOne leader accepts writesSimple; leader is a write bottleneck & SPOF (failover needed)
Multi-leaderSeveral leaders (per region)Low write latency per region; write conflicts to resolve
LeaderlessAny replica; quorumHighly 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.

StrategyHowDownside
RangeBy key ranges (A–M, N–Z)Hot spots on sequential keys
HashBy hash(key)Even spread; loses range queries
DirectoryLookup table maps key→shardFlexible; 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.

text
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.

LeanSystems
CP (consistency)HBase, MongoDB (majority), Spanner, etcd, ZooKeeper
AP (availability)Cassandra, DynamoDB, Riak (tunable)

Isolation Levels & Anomalies

LevelDirty readNon-repeatablePhantom
Read UncommittedPossiblePossiblePossible
Read CommittedNoPossiblePossible
Repeatable ReadNoNoPossible*
SerializableNoNoNo

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 CommitSaga
ModelCoordinator: prepare then commitSequence of local txns + compensations
AtomicityStrong (all or none)Eventual; roll back via compensations
BlockingLocks held; blocks if coordinator diesNon-blocking, per-step commit
FitFew nodes, short txnsMicroservices, 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.

text
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.

Section navigation