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.
Key points
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.
Single copy. No spare.
Single copy. No replica.
Serving. One failure away from total outage. That is why (1), (2) and (3) all break here.
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.
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.
replicated
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.
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_id | user_id % 4 | Shard |
|---|---|---|
| 0, 4, 8, 12 | 0 | Shard 0 |
| 1, 5, 9, 13 | 1 | Shard 1 |
| 2, 6, 10, 14 | 2 | Shard 2 |
| 3, 7, 11, 15 | 3 | Shard 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.
Even split. user_id % 4 spreads 0 to 15 as 4 users per shard. A bad key would pile users on one shard.
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.
Hotspot. Shard 2 carries Katy Perry plus Justin Bieber plus Lady Gaga reads and overloads while other shards sit idle.
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.
This runnable Python models the two mechanisms that matter. Lines that matter are marked.
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.
return user_id % NUM_SHARDS is the whole routing contract. A bad sharding key here creates uneven shards and future resharding.queue_for_photo_processing.append(job) enables independent scaling. Add workers when the queue grows. Remove workers when it stays empty.popleft() shows durability under crash. If workers stop, jobs wait. If producers stop, workers drain remaining jobs.| Who / what | Why |
|---|---|
| Netflix: asynchronous multi data center replication | Why: 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 visitors | Why: delay sharding complexity while vertical scaling still fits. |
| Photo processing workers: message queue between web servers and workers | Why: cropping and similar tasks are slow, so buffer requests and scale workers alone. |
| NoSQL store for non relational functionalities | Why: reduce load on sharded relational databases. |
| CDN for www.mysite.com and api.mysite.com static assets | Why: serve users from edge and reduce origin and database load. |
| GeoDNS for US-East and US-West routing | Why: send users to the closest data center in normal operation and reroute on outage. |
| Logging plus metrics plus automation tools in every data center | Why: find errors fast and keep deploys consistent, which is maintainability at scale. |
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.