All lessonsDay 15A-Foundations2026-10-10

Reliable, scalable, maintainable: the three pillars

Your photo site runs in US-East and US-West. One night US-West goes offline completely. Can US-East absorb 100% of traffic and still show every user their data.

x / 100-xTraffic split: x% in US-East and (100 minus x)% in US-West. Failover: 100% to US-East.
24 TB RAMVertical limit example: one database server with 24 TB of RAM.
10M+ / 1 masterScale without sharding example: stackoverflow.com in 2013 had over 10 million monthly unique visitors with 1 master database.
user_id % 4Shard routing example: user_id % 4 into Shard 0, 1, 2, 3.

Key points

  • Reliability comes from redundancy across tiers and data centers.
  • Scalability comes from decoupling tiers so each tier grows alone.
  • Maintainability comes from logging, metrics and automation that work the same in every data center.
O1Derive which shard stores a given user_id using modulo routing.
O2Break down what fails during a full data center outage and where requests go.
O3Argue when to choose vertical scaling, sharding, or a queue under cost and failure tradeoffs.
REQ

Requirements


  1. (1)Reliable: the system stays up when a server or a whole data center fails.
  2. (2)Scalable: each part handles more users and data without blocking other parts.
  3. (3)Maintainable: the team can test, deploy, find errors and understand health at scale.
01

Problem: one place for everything breaks all three requirements


One web server plus one database is simple to run. Growth breaks it in three ways.

First, a single server is a single point of failure. One failure stops service. That breaks (1).

Second, web, cache and database grow at different rates. Scaling them together wastes resources. That breaks (2).

Third, manual deploys and per server logs hide errors. The team cannot see database tier health or cache tier health. That breaks (3).

The rest of this lesson adds one mechanism per break.

Try it: press Fail the server. With one box per tier there is nowhere for traffic to go.

status: serving
Web server
UP

Single copy. No spare.

Database
UP

Single copy. No replica.

Serving. One failure away from total outage. That is why (1), (2) and (3) all break here.

One web server plus one database: simple, but any single failure stops service and every tier must scale together.
02

Mental model: stateless tiers behind routing


Picture this flow.

User to DNS to CDN to load balancer to web servers to cache to database. Each box can have many copies.

The model has three rules from the chapter summary.

Keep web tier stateless. Any web server can serve any request.

Build redundancy at every tier. No tier depends on one machine.

Cache data as much as you can. Host static assets in CDN.

This picture lets us reason about failure. If one web server dies, the load balancer sends traffic elsewhere. If one data center dies, DNS routing sends traffic to the survivor. If the database fills, we split only the data tier.

Walk it: press Next to follow one photo view from user to data. Each lit box is the current holder of the request.

DNSgeo route to closest DC
CDNstatic assets at edge
Load balancerspread to web pool
Web serversstateless, any can serve
Cacheas much as you can
Databasesplit only this tier later
  1. User asks for a photo. DNS answers with the closest healthy data center IP.
  2. CDN serves static bytes for www.mysite.com and api.mysite.com without touching origin.
  3. Load balancer picks a healthy web server. Dead servers get no traffic.
  4. Stateless web server handles any request, so any survivor can take over.
  5. Cache absorbs hot reads before they reach the database.
  6. Database answers only cache misses. When it fills, only this tier is sharded.
Stateless web tier plus redundancy at every tier plus CDN caching. Failure at one copy reroutes to another copy.
03

Mechanism: multiple data centers with geo routing and failover


A geoDNS service resolves a domain to an IP address based on user location. Failover means directing all traffic to a healthy data center after an outage.

Normal operation uses two data centers. Users are geo routed to the closest data center. The split is x% in US-East and (100 minus x)% in US-West.

Failover operation uses one data center. In the example, data center 2 in US-West is offline. Then 100% of traffic is routed to data center 1 in US-East.

Three challenges must be solved.

Traffic redirection. GeoDNS must direct traffic to the correct data center.

Data synchronization. Users in different regions may use different local databases or caches. Traffic routed to a new data center may find data unavailable. A common strategy is to replicate data across multiple data centers.

Test and deployment. The site must be tested at different locations. Automated deployment tools keep services consistent through all the data centers. That supports (3).

Failure case. US-West goes offline without cross data center replication. US-East takes 100% of traffic but misses US-West writes. That breaks (1). The fix is asynchronous multi data center replication. Its price is complexity and delayed copies.

Drive it: set the normal split with x, then kill US-West. Toggle replication to see when 100% traffic still means missing data.

60%
US-East load
60%
US-West load
40%
Data visible in East
100%

replicated

GeoDNS moves traffic, replication moves data. Traffic failover without data replication still breaks (1).
04

Mechanism: message queue to decouple and absorb spikes


A message queue is a durable component stored in memory. It supports asynchronous communication. It serves as a buffer and distributes asynchronous requests.

A producer creates messages and publishes them to the queue. A consumer connects to the queue and performs actions defined by the messages. The consumer can subscribe for push or consume by pull.

Decoupling helps (1) and (2). The producer can post a message when the consumer is unavailable. The consumer can read messages when the producer is unavailable.

Concrete use case from the chapter. The application supports photo customization including cropping, sharpening and blurring. Those tasks take time to complete. Web servers publish photo processing jobs to the queue. Photo processing workers pick up jobs and perform customization asynchronously.

Scaling rule. The producer and the consumer can be scaled independently. When the size of the queue becomes large, add more workers to reduce processing time. If the queue is empty most of the time, reduce the number of workers.

Failure case. Workers crash. Messages stay durable in the queue. Web servers keep accepting uploads. Latency grows but work is not lost. That protects (1). The price is async behavior. Users do not get instant results.

Drive it: raise arrivals above workers to grow the queue, then add workers to drain it. Crash workers to see durability without loss.

6
4
Queue depth: 0Processed: 0Lost: 0

Producer rate and consumer rate are independent. Depth signals when to add or remove workers. Crashes add delay, not loss.
05

Mechanism: database scaling by vertical growth then sharding


There are two broad approaches.

Vertical scaling means adding more power to an existing machine. Power means CPU and RAM and DISK. Example limit cited is 24 TB of RAM on Amazon Relational Database Service. Stackoverflow.com in 2013 handled over 10 million monthly unique visitors with 1 master database.

Vertical scaling breaks (1) and (2) at scale. Hardware has limits. A single server is not enough for a large user base. Risk of single point of failures is greater. Powerful servers are much more expensive.

Horizontal scaling means adding more servers. In databases this is called sharding. A shard is a smaller part of a large database. Each shard shares the same schema. The actual data on each shard is unique to the shard.

Routing uses a sharding key. The sharding key is one or more columns that determine how data is distributed. In the example the key is user_id. The hash function is user_id % 4.

Result mapping:

user_iduser_id % 4Shard
0, 4, 8, 120Shard 0
1, 5, 9, 131Shard 1
2, 6, 10, 142Shard 2
3, 7, 11, 153Shard 3

A sharding key allows routing of queries to the correct database. The most important criterion is a key that can evenly distribute data.

Drive it: type any user_id to route it, then check the 0 to 15 grid against the table above.

user 13 -> shard 1

Even split. user_id % 4 spreads 0 to 15 as 4 users per shard. A bad key would pile users on one shard.

One line, user_id % 4, decides correctness. Change it and data must move.
06

Weakness stated honestly: sharding adds three new problems


Sharding helps (2) but hurts (1) and (3) unless handled.

Resharding data. Resharding means updating the sharding function and moving data around. It is needed when a single shard can no longer hold more data due to rapid growth. It is also needed when certain shards experience exhaustion faster due to uneven distribution. Consistent hashing is named as a commonly used technique for this problem.

Celebrity problem. This is also called a hotspot key problem. Excessive access to a specific shard can cause server overload. Example given is data for Katy Perry and Justin Bieber and Lady Gaga ending up on the same shard. For social applications that shard will be overwhelmed with read operations. A fix is to allocate a shard for each celebrity. Each shard might even require further partition. The price is special cases and more shards to manage.

Join and de normalization. Once a database has been sharded across multiple servers, it is hard to perform join operations across shards. A common workaround is denormalization. Denormalization stores data so queries can be performed in a single table without cross shard joins. The price is duplicated data and harder updates.

Drive it: put all three celebrities on Shard 2, watch overload, then split them one per shard.

Shard 0 reads
10%
Shard 1 reads
10%
Shard 2 reads
90%
Shard 3 reads
10%

Hotspot. Shard 2 carries Katy Perry plus Justin Bieber plus Lady Gaga reads and overloads while other shards sit idle.

Even data counts can still mean uneven read heat. Splitting hot keys trades special case logic for headroom.

A related fix in the design is to move some non relational functionalities to a NoSQL data store. That reduces database load. The updated design also adds logging and metrics and monitoring and automation tools.

Maintainability tooling is required for (3).

Logging. Monitor error logs to identify errors and problems. Monitor at per server level or aggregate to a centralized service for search and viewing.

Metrics. Collect host level metrics such as CPU and Memory and disk I/O. Collect aggregated level metrics such as performance of the entire database tier and cache tier. Collect key business metrics such as daily active users and retention and revenue.

Automation. Build or leverage automation tools to improve productivity. Continuous integration means each code check in is verified through automation. Automate build and test and deploy to improve developer productivity.

07

Code: shard routing plus queue buffering


This runnable Python models the two mechanisms that matter. Lines that matter are marked.

Show code
from collections import deque

# Shard routing matters because it decides correctness.
# Chapter function: user_id % 4
NUM_SHARDS = 4

def get_shard(user_id: int) -> int:
    # This one line is the sharding function.
    # Change this and you must move data. That is resharding.
    return user_id % NUM_SHARDS

# Message queue matters because it decouples rate of arrival from rate of work.
queue_for_photo_processing = deque()

def publish(job: str) -> None:
    # Producer does not wait for workers. This protects availability.
    queue_for_photo_processing.append(job)

def consume_one() -> str | None:
    # Consumer works even if producer is offline, as long as queue has jobs.
    if queue_for_photo_processing:
        return queue_for_photo_processing.popleft()
    return None

if __name__ == "__main__":
    # Illustrative user_ids. Routing must match Figure 1-22 pattern.
    for uid in [0, 1, 2, 3, 4, 5, 13, 14, 15]:
        print(f"user {uid} -> shard {get_shard(uid)}")

    publish("crop:photo42")
    publish("blur:photo43")
    print(consume_one())
    print(consume_one())
    print(consume_one())

Expected output

user 0 -> shard 0
user 1 -> shard 1
user 2 -> shard 2
user 3 -> shard 3
user 4 -> shard 0
user 5 -> shard 1
user 13 -> shard 1
user 14 -> shard 2
user 15 -> shard 3
crop:photo42
blur:photo43
None

Notes on the lines that matter.

08

Field guide: who uses what in this chapter, and why


Who / whatWhy
Netflix: asynchronous multi data center replicationWhy: keep data available in another region so failover can serve 100% traffic after an outage.
Stack Overflow 2013: 1 master database for over 10 million monthly unique visitorsWhy: delay sharding complexity while vertical scaling still fits.
Photo processing workers: message queue between web servers and workersWhy: cropping and similar tasks are slow, so buffer requests and scale workers alone.
NoSQL store for non relational functionalitiesWhy: reduce load on sharded relational databases.
CDN for www.mysite.com and api.mysite.com static assetsWhy: serve users from edge and reduce origin and database load.
GeoDNS for US-East and US-West routingWhy: send users to the closest data center in normal operation and reroute on outage.
Logging plus metrics plus automation tools in every data centerWhy: find errors fast and keep deploys consistent, which is maintainability at scale.
09

Interview ears: follow ups with answers


10

Judgment quiz


Q1Predict. US-West is offline. GeoDNS sends 100% to US-East. US-East has no replica of recent US-West photo jobs. Which requirement breaks and what does the user see.

Uploads may succeed but recent jobs or reads appear missing until replication or queue drain restores them. Traffic failover without data replication still breaks reliability.

Q2Compare. Vertical scaling to 24 TB RAM versus sharding into 4 shards by user_id % 4. Which improves (2) with lower operational cost at 10 million monthly uniques, and when does the answer flip.

At Stack Overflow 2013 scale one master handled over 10 million monthly uniques. Past hardware limits or single point of failure risk, sharding wins for (2) despite resharding and join costs.

Q3Choose under tradeoff. Your queue grows every afternoon because photo workers are slow. You can add workers, denormalize the job table, or move non relational job metadata to NoSQL. You have one week and oncall pain is high. What do you choose and what price do you accept.

Worker count is the independent scaling knob for a growing queue. Schema changes wait until the database, not worker speed, is proven to be the limit.

Sources: Designing Data-Intensive Applications Chapter 1 (reliability, scalability, maintainability); System Design Interview Chapter 1 scale from zero to millions (DNS, CDN, load balancer, queue, sharding, geoDNS, failover, logging, metrics, automation).