Practice this topic in a realistic system design interview
An application that works well for a small group of users can struggle once that group grows. More people sign up, more requests arrive every second, and more data needs to be stored. At some point, the original setup cannot keep up.
Scalability is a system's ability to handle more load by adding resources. This chapter covers what scalability means, the different ways to scale a system and their trade-offs, and how an application can grow one step at a time.
Suppose we have built a social media application where users create profiles, publish posts, and read their feeds. At first, the application and the database run on one server, and everything works well.
But as more people use it, more requests arrive and more data needs to be stored. Eventually, the server cannot finish the work fast enough. Requests start waiting, pages load slowly, and some requests time out.
We need a way to handle that extra load while keeping the application responsive. That is the problem scalability solves.
Scalability is a system's ability to handle more load by adding resources.
Suppose our application handles 1,000 requests per second, and we want it to handle 5,000 while still meeting our response-time requirements. We might give the server more memory, or we might add servers to share the work.
Doubling the resources will not always double the capacity. A second application server does not help if every request waits on the same slow database. What matters is whether the added resources help the system handle more of the work we need it to do.
Scalability is not a yes-or-no property. A system can scale well for one kind of growth and poorly for another. It helps to be specific about what is growing: the number of requests, the number of active users, or the amount of stored data.
Before choosing an approach, we look at what is increasing. Three measurements come up often.
These tell us different things. One million registered users does not mean one million people are using the application at once. Most of them are idle at any given moment, and only a fraction are sending requests.
The type of request matters too. Reading a user profile takes much less work than processing a video. A thousand profile reads per second and a thousand video uploads per second are very different loads, even though the request count is the same.
So user count alone is not enough. We need to know how often requests arrive and what each request asks the system to do.
Next, we look at how the system behaves under that load. Suppose requests normally take 50 milliseconds, but now they take 500 milliseconds, and some fail.
We check CPU usage, memory, disk activity, and database queries to find where the time goes. Each of these points to a different cause.
Maybe the application is using all its CPU. Or maybe it is waiting for a slow query that we could improve with an index.
The component that limits how much work the system can handle is called the bottleneck. When we remove one bottleneck, another component usually becomes the next limit, so this is a loop rather than a one-time step.
Adding application servers will not fix a slow database query. So we identify the cause before adding resources. The rest of this chapter follows that loop: find what is limiting the system, change that part, and check the result.
Let's start with the simplest approach to scale: make the existing machine more powerful. This is called vertical scaling, or scaling up.
Loading simulation...
Suppose our server has 4 CPU cores and 16 GB of memory. We could move the application to a machine with 8 cores and 32 GB of memory.
If CPU or memory was the limit, that upgrade may let us handle more requests. The application can often keep working in the same way, with few or no code changes.
But the upgrade needs to match the problem. More CPU will not do much if requests are waiting on a slow disk. In that case, faster storage may be the more useful change.
The same logic applies to the other resources. A full memory calls for more RAM, and a saturated network link calls for a faster network interface. The upgrade that helps is the one that matches the resource that is actually full.
As our application grows further, the application and the database begin competing for the upgraded server's resources. Both want CPU, memory, and disk at the same time, and each slows the other down.
We move the database to its own machine so each can be sized separately.
The database can get more memory, since databases use memory to keep frequently accessed data close and avoid disk reads. The application server gets CPU for processing requests. Each machine now runs one kind of workload, and we can upgrade one without touching the other.
This is still vertical scaling. We have two machines instead of one, but each is scaled up on its own.
Later, the application server reaches its own limit. Instead of upgrading it again, let's look at how we can distribute its work across multiple machines.
The second approach is horizontal scaling, also called scaling out. We run the application on several servers and put a load balancer in front of them.
Loading simulation...
The load balancer distributes requests so each server handles part of the traffic.
When we need more capacity, we add another server. Auto-scaling can adjust that number as demand changes, adding servers when traffic rises and removing them when it drops.
If one server becomes unhealthy, the load balancer stops sending it requests. The others can continue serving traffic, provided they have enough spare capacity to absorb its share.
But before we can freely send requests to any server, there is one detail we need to handle.
Suppose a user logs in through Server 1, which keeps their session only in its own memory. Their next request goes to Server 2. Server 2 does not have that session, so it may not recognize them as logged in.
We could keep sending this user back to Server 1. That is called a sticky session. The load balancer remembers which server each user landed on and routes them there every time.
But sticky sessions give us less freedom to distribute traffic, since some servers end up carrying more active users than others. And if Server 1 fails, its session data is lost, and every user on it has to log in again.
We want any healthy application server to be able to handle the next request.
We can solve this by moving session data into a shared store, such as Redis. Now, when a request arrives, whichever application server receives it can look up the same session.
We also put files, such as user uploads, in shared storage like S3, so they do not depend on one application server's local disk.
With this setup, the application servers are stateless. The system still has state. That state now lives outside the individual application servers.
This makes it much easier to add a server, remove one, or replace one after a failure. A new server has nothing to load before it can serve requests, and a failed server takes nothing with it.
We now have both approaches in use, so this is a good point to compare them side by side.
So vertical scaling is simpler, while horizontal scaling gives us more room to grow at the cost of more machines, network communication, and shared data to manage.
We can also combine them. Our application now runs on several servers, while the database can still run on one larger machine.
The next question is how to keep that database from becoming overwhelmed as those application servers send it more queries.
Suppose many of those queries read the same popular user profiles. Every request for a popular profile runs the same query and gets the same result. We can avoid repeating that work by keeping copies in a cache.
The application checks the cache first. If the profile is there, it returns that value. Otherwise, it reads from the database and puts the result in the cache for later requests. This is called the cache-aside pattern.
It reduces database work and can make repeated reads faster, since reading from memory takes less time than running a database query.
When a profile changes, though, we need to update or remove its cached copy, or let it expire. Otherwise, users may keep seeing the old value.
There are three common ways to handle that. We can delete the cached entry when the profile changes, so the next read refills it from the database. We can write the new value to both the database and the cache, which keeps the cache warm but gives the application two places to keep in sync. Or we can give each entry a time to live and accept that users may see the old value until it expires.
As the cache grows, one cache server may run out of memory or processing capacity. We can spread entries across several nodes, using each key to decide where an entry belongs.
That gives us more total space, and it lets different nodes handle different requests at the same time.
Caching helps with repeated reads, but requests for uncached data still reach our database. A feed that is assembled fresh, a profile that is rarely viewed, or a search over recent posts all need a real query.
If those reads become too much for the database, we can add read replicas. A replica is a separate database that keeps a copy of the primary's data.
Writes go to the primary, which sends its changes to the replicas. Our application can then spread profile and post reads across those copies.
Every write still goes through the primary, so replicas increase read capacity without changing the primary's write limit.
There is a small delay between a write on the primary and the moment it appears on a replica. A user who has just edited their profile might read it back from a replica that has not caught up yet. Most reads tolerate that delay. For the few that must reflect the latest write, the application can read from the primary instead.
Later, suppose writes or data size outgrow the primary database. Replicas do not help here, because every replica stores the full data set and every write still lands on the primary.
We can consider sharding, which means dividing data across separate databases.
For our application, we could use a user ID to decide which database stores a user's profile and posts. That user ID is the shard key.
Each database, or shard, handles only its portion of the data. Different shards can process writes at the same time, spreading both storage and write traffic across machines.
Sharding does add work for the application. A query that touches one user, such as loading their profile, goes to one shard. A query that spans many users, such as a search across all posts, has to run on every shard and merge the results. Choosing a shard key that keeps common queries inside one shard makes this much easier.
Some work does not need to finish before we respond to the user. For example, after a user publishes a post, we may need to send notifications to their followers. A user with many followers would wait a long time if we sent every notification before confirming the post.
We can save the post in the database, place a notification job on a queue, and respond without waiting for every notification to be sent. Background workers take jobs from the queue and process them separately.
During a short traffic spike, jobs can wait in the queue while workers catch up. If the backlog keeps growing, we may need more workers, provided the work can be shared safely between them.
The queue has not reduced the total work. It has allowed us to process that work separately from the user's request, so the user's response time no longer depends on it.
After a change, we test with a workload similar to what we expect. We compare how many requests finish, how long they take, and how many fail.
We look beyond the average response time, because some users may still experience long waits even when the average looks good. An average of 60 milliseconds can hide a small group of requests that take two seconds. Percentiles such as p95 and p99 show what those slower requests look like.
We also check cost, since the new setup needs to be affordable at the load we expect. And we make sure data is still handled correctly. A faster system that loses writes or returns the wrong profile has not improved.
There is no fixed user count for each step we have discussed. One application may need sharding early because each user stores a large amount of data. Another may run on a single database for years because each user reads a little and writes rarely. The right change depends on the work our application does and whether the result meets its requirements.
The main idea is to find the bottleneck, make a change that addresses it, and measure the result.
In this chapter, we scaled our social media application step by step. Each step addressed the limit the previous step exposed.
Put together, the application now looks like this.
That is the core of scalability: identifying what breaks next and scaling that part of the system.
As an exercise, pick a real-world application and think about how you would scale it from a single request to millions of requests. Which component would reach its limit first, what would you change, and what would you measure to know it worked?
10 quizzes