Skip to content
LowLevelDesign Mastery

Scaling Databases

Making databases handle more load

As your application grows, your database becomes a bottleneck. More users mean more queries, and a single database server can only handle so much.

Database scaling journey: a single database at 1000 QPS gets overloaded at 10,000 QPS, then read replicas and pooling reach 100,000 QPS

Vertical scaling means making your existing server bigger—more CPU, more RAM, faster disks.

Vertical scaling: upgrading from 4 cores and 16GB to 16 cores and 64GB raises capacity to 4000 QPS but hits cost and hardware limits

Pros:

  • Simple (just upgrade hardware)
  • No code changes needed
  • Works immediately

Cons:

  • Expensive (bigger servers cost more)
  • Hardware limits (can’t scale infinitely)
  • Single point of failure

Horizontal scaling means adding more servers to share the load.

Horizontal scaling: a single database scales out to a primary handling writes that replicates to two read replicas for 3000 QPS

Pros:

  • Can scale infinitely (add more servers)
  • Cost-effective (cheaper servers)
  • High availability (if one fails, others continue)

Cons:

  • More complex (need replication, load balancing)
  • Requires code changes
  • Data consistency challenges

Major companies use different scaling strategies based on their needs:

The Challenge: A growing SaaS application needs more database capacity but wants to keep architecture simple.

The Solution: Many startups use scale-up initially:

  • Start: Single database server (4 CPU, 16GB RAM)
  • Growth: Upgrade to larger server (16 CPU, 64GB RAM)
  • Benefit: Simple, no code changes, immediate improvement

Example: A B2B SaaS application:

  • Year 1: 1,000 users → Single server handles load
  • Year 2: 10,000 users → Upgrade to 4x larger server
  • Year 3: 50,000 users → Upgrade to 16x larger server

Impact: Handled 50x growth with scale-up. Simple operations. Cost-effective for moderate scale.

Limitation: Eventually hits hardware limits. At 100,000+ users, need scale-out.

The Challenge: Large applications need to scale beyond single server limits. Need high availability and fault tolerance.

The Solution: Major platforms use scale-out:

  • Facebook: Thousands of MySQL servers, read replicas, sharding
  • Amazon: Distributed databases, read replicas across regions
  • Netflix: Multiple database clusters, read replicas for global access

Example: E-commerce platform scaling:

  • Start: Single database (handles 1,000 orders/day)
  • Growth: Add read replicas (handles 10,000 orders/day)
  • Scale: Shard database (handles 1,000,000 orders/day)

Impact: Scales to millions of users. High availability (survives server failures). Global distribution.

The Challenge: Instagram serves billions of photo views daily. Most operations are reads (viewing photos, profiles, feeds).

The Solution: Instagram uses extensive read replicas:

  • Primary: Handles writes (photo uploads, comments, likes)
  • Replicas: Handle reads (photo views, profile views, feed generation)
  • Scale: 100+ read replicas worldwide

Why Read Replicas?

  • 99% of operations are reads
  • Read replicas distribute load
  • Primary focuses on writes

Impact: Handles billions of reads daily. Primary not overloaded. Fast read performance globally.

The Challenge: APIs handling thousands of requests per second. Creating database connections is expensive (50ms per connection).

The Solution: Major APIs use connection pooling:

  • GitHub: Connection pool of 100 connections, reused across requests
  • Stripe: Connection pool per API server, handles millions of requests
  • Twitter: Connection pooling essential for high throughput

Example: API handling 10,000 requests/second:

  • Without pool: 10,000 connections × 50ms = 500 seconds overhead
  • With pool: 100 connections reused, 5ms per request = 50 seconds overhead
  • 10x improvement

Impact: Reduced connection overhead by 90%. Handles 10x more requests. Lower latency.

Read-Write Splitting: E-commerce Platforms

Section titled “Read-Write Splitting: E-commerce Platforms”

The Challenge: E-commerce platforms have read-heavy workloads (product browsing) and write operations (orders, inventory updates).

The Solution: E-commerce platforms use read-write splitting:

  • Writes: Go to primary (orders, inventory updates)
  • Reads: Go to replicas (product catalog, search, recommendations)

Example: Product page load:

  • Read: Product data from replica (fast, doesn’t block writes)
  • Write: Order creation goes to primary (ensures consistency)

Impact: Product pages load faster. Order processing not blocked by reads. Better user experience.



Read replicas are copies of the primary database that handle read operations. They replicate data from the primary and distribute read load.

Read replicas: the primary database streams changes asynchronously to two eventually consistent replicas that serve reads

Key Characteristics:

  • Primary handles all writes
  • Replicas handle reads
  • Replication is asynchronous (eventual consistency)
  • Load distribution across multiple servers

Read-write splitting routes writes to the primary and reads to replicas.

Read-write splitting: a router sends INSERT, UPDATE and DELETE queries to the primary and SELECT queries to a replica

Benefits:

  • Distributes read load
  • Primary focuses on writes
  • Better overall performance

Replication lag is the delay between a write on the primary and when it appears on replicas.

Replication lag: a write to the primary takes 100-500ms to reach the replica, so a read 50ms later returns stale data

Problem: Users might read stale data immediately after writing.

Solution: Read-after-write consistency—read from primary for a short time after writing.


Connection pooling maintains a cache of database connections that can be reused, reducing connection overhead.

Connection overhead without a pool at 50ms per query versus borrowing and returning pooled connections at 5ms per query

Connection overhead:

  • Creating connection: ~20ms
  • Authentication: ~10ms
  • Closing connection: ~5ms
  • Total overhead: ~35ms per query

With pooling: Reuse existing connections, overhead: ~1ms


Connection pool: a request borrows one of several active or idle connections from the pool and returns it after use

Pool Management:

  • Min pool size: Keep minimum connections ready
  • Max pool size: Maximum concurrent connections
  • Idle timeout: Close idle connections after timeout
  • Connection timeout: Wait time if pool is full

How database scaling affects your class design:

💡 Tip: Click dropdown to switch between languages
Connection Pool
from queue import Queue
from threading import Lock
import psycopg2
from contextlib import contextmanager
class ConnectionPool:
def __init__(self, connection_string, min_size=5, max_size=20):
self.connection_string = connection_string
self.min_size = min_size
self.max_size = max_size
self.pool = Queue(maxsize=max_size)
self.lock = Lock()
self.active_connections = 0
# Initialize pool with min connections
for _ in range(min_size):
conn = psycopg2.connect(connection_string)
self.pool.put(conn)
@contextmanager
def get_connection(self):
conn = None
try:
# Get connection from pool or create new one
if not self.pool.empty():
conn = self.pool.get_nowait()
elif self.active_connections < self.max_size:
conn = psycopg2.connect(self.connection_string)
self.active_connections += 1
else:
# Wait for available connection
conn = self.pool.get(timeout=5)
yield conn
finally:
# Return connection to pool
if conn:
self.pool.put(conn)
def close_all(self):
while not self.pool.empty():
conn = self.pool.get()
conn.close()
💡 Tip: Click dropdown to switch between languages
Read-Write Splitting
class ReadWriteRouter:
def __init__(self, primary_pool, replica_pools):
self.primary_pool = primary_pool
self.replica_pools = replica_pools
self.replica_index = 0
def execute(self, query, params=None, is_write=False):
# Route writes to primary, reads to replicas
if is_write or self._is_write_query(query):
pool = self.primary_pool
else:
# Round-robin across replicas
pool = self.replica_pools[self.replica_index]
self.replica_index = (self.replica_index + 1) % len(self.replica_pools)
with pool.get_connection() as conn:
cursor = conn.cursor()
cursor.execute(query, params)
if is_write or self._is_write_query(query):
conn.commit()
return cursor.rowcount
else:
return cursor.fetchall()
def _is_write_query(self, query):
query_upper = query.strip().upper()
return query_upper.startswith(('INSERT', 'UPDATE', 'DELETE', 'CREATE', 'DROP', 'ALTER'))


Now that you understand scaling, let’s explore sharding and partitioning to distribute data across multiple databases:

Next up: Sharding & Partitioning — Learn horizontal vs vertical partitioning and how to implement shard-aware repositories.