AlgoMaster Logo

Scalability

High Priority11 min readUpdated September 15, 2026
AI Mock Interview

Practice this topic in a realistic system design interview

Listen to this chapter
Unlock Audio

Scalability

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.

1. The Problem

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.

2. What Scalability Means

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.

3. Measuring the Load

Before choosing an approach, we look at what is increasing. Three measurements come up often.

Scroll
MetricWhat It MeasuresExample
Requests per secondHow often requests arrive1,000 requests/s
Concurrent usersUsers active at the same time20,000 users online now
Data volumeAmount of data stored2 TB of profiles, posts, and images

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.

Scroll
RequestWork InvolvedRelative Cost
Read a profileOne indexed lookup, small responseLow
Load a feedMany lookups, sorting, assembling a pageMedium
Upload a videoTranscoding, large file storage, heavy CPU useHigh

So user count alone is not enough. We need to know how often requests arrive and what each request asks the system to do.

4. Finding the Bottleneck

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.

Scroll
What We CheckSymptomLikely Cause
CPU usageSitting at 100%The application is doing more computation than the cores can handle
MemorySwapping or out-of-memory errorsThe working data no longer fits in RAM
Disk activityHigh I/O waitThe database is reading from slow storage
Database queriesA few queries take most of the timeA missing index or an inefficient query

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.

5. Vertical Scaling (Scale Up)

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.

Trade-offs

Pros

  • Simplicity: We still have one machine to run, and we do not need to change how requests are handled.
  • Few or no code changes: The application usually works the same way on the bigger machine.
  • No network between components: Everything stays local, so there is no added latency from machines talking to each other.

Cons

  • Hardware ceiling: There is a limit to how powerful one machine can become.
  • Cost: Larger machines get expensive, and the price often grows faster than the capacity.
  • Downtime: An upgrade may require a restart or some downtime while the application moves to the new machine.
  • Single point of failure: If this is our only server, there is still no other server to take over when it fails.

6. Separating the Database

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.

7. Horizontal Scaling (Scale Out)

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.

The Session Problem

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.

Stateless Application Servers

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.

8. Vertical vs Horizontal Scaling

We now have both approaches in use, so this is a good point to compare them side by side.

Scroll
Vertical ScalingHorizontal Scaling
How it worksMake one machine more powerfulAdd more machines behind a load balancer
Code changesUsually noneServers must be stateless
LimitThe largest available machineAdd servers as long as shared components keep up
When a server failsNothing takes overThe load balancer routes around it
Cost curveGrows steeply at the high endRoughly linear per server
Operational complexityLowHigher: more machines, network calls, shared data

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.

9. Caching

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.

Scaling the Cache

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.

10. Read Replicas

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.

11. Sharding

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.

12. Queues and Background Work

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.

13. Verifying the Change

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.

Scroll
MeasurementBeforeAfterWhat It Tells Us
Requests completed per second1,0005,000Whether capacity actually grew
Median response time50 ms45 msWhat a typical request experiences
99th percentile response time500 ms120 msWhat the slowest 1% of requests experience
Failed requests2%0.1%Whether the system is dropping work under load

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.

Summary

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.

Scroll
StepChangeBottleneck It Addressed
1Upgrade the machineCPU or memory on a single server
2Separate the databaseApplication and database competing for one machine
3Distribute traffic across stateless serversOne application server at its limit
4Add cachingRepeated reads of the same data
5Add read replicasTotal read volume on the database
6Shard the databaseWrite volume and data size on the primary
7Use queues for background workSlow work sitting inside the user's request

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?

Quiz

Scalability Quiz

10 quizzes