System design for heavy load
When one server is not enough, we add more servers instead of a bigger server, and we spread the load between them. The database is the hardest part. With Replicas we spread the reads, and with Sharding we split the data itself. Each of these has a cost that we must know.
Author: bezzad
The problem: the 12 o’clock flash sale
At 12 noon, our shop sells one thousand phones at a discount. In the first minute, about two hundred thousand people come in and click the buy button many times. We have three conditions:
- No more than one thousand phones are sold.
- The system does not go down.
- Each user buys only one phone.
Today, everything runs on one server and one database. This server stops working in the first few seconds.
First step: calculate the real numbers
Before you design, do a quick rough calculation. This also earns points in an interview.
- Two hundred thousand users in sixty seconds means about 3300 users per second.
- If each user clicks or refreshes five times on average, that means about 16 thousand requests per second.
- But only one thousand people really buy something. So more than 99% of requests get the answer “sold out” in the end.
Important result: not all requests need to reach the database. We must answer most of them cheaply and quickly.
Two ways to grow
- The vertical way (Scale Up). Buy a stronger server. The code does not change. But it has a limit, it is expensive, and if that one server goes down, everything goes down.
- The horizontal way (Scale Out). Put several normal servers side by side and spread the load. It has almost no limit, and the failure of one server does not matter. But it has one condition: the app must be stateless.
What does stateless mean? It means no user data stays in the memory of one specific server:
- The user’s first request goes to server 1, and the shopping cart is saved in the memory of server 1.
- The same user’s second request goes to server 2.
- Server 2 does not see the cart. The user thinks the cart is now empty.
- So keep user data in a shared place, like Redis or the database.
Spreading the load (Load Balancing)
A Load Balancer sits in front of the servers and gives each request to one of them. Some common methods:
- The turn-based method (Round Robin). In order: first, second, third, then first again. It is simple and works well for requests of the same size.
- The least connections method (Least Connections). Give it to the server that has the least work right now. It is better for requests whose time is very different.
- The key-based method (Hash). Requests from one user always go to the same server. It is only needed when the server has state. Usually it is better to remove that state.
Another important job of the Load Balancer is the health check. Every few seconds, it asks each server “are you healthy?”. It takes an unhealthy server out of the list. But between the server failing and the next check, some requests reach the broken server.
Live example
Send some requests. Then break server 2 and send again:
Designing the flash sale
Now let us put everything together:
- Reduce the load before the server. Put the product page and images on a CDN. Set a request limit for each user and each IP. Repeated clicks do not reach the servers.
- Stateless API servers, many of them. Behind the Load Balancer, with horizontal Scale.
- Decrease the stock atomically. This is the most important part (code below). Reading the stock and decreasing it must be one action, not two separate actions.
- Give heavy work to a queue. After a successful reservation, only send one message to the queue and tell the user “reserved”. A Worker creates the order, the invoice and the email at its own speed. The queue smooths out the peak load.
- Protect the database. Send the reads to the Replicas. Writes go only to the Primary.
Database: Replica
Most systems read much more than they write. A customer looks at a product a hundred times and buys once. So:
- One Primary takes all the writes.
- Several Replicas have a copy of the data, and the reads are spread between them.
- The copy usually has a delay. From a few milliseconds to a few seconds. This is called Replication Lag.
Let us see the delay problem with an example: a customer places an order and right away goes to the “My orders” page. This page reads from a Replica. The order has not reached the Replica yet. The customer thinks the order is lost.
Solution: for reads that come right after the same user’s write, read from the Primary. This is called reading your own writes (Read Your Writes).
Database: Sharding
When the size of the data or the number of writes is too much for one server, we split the data itself. We call each piece a Shard, and it is on a separate database.
The most important decision is choosing the Shard key:
- The key must spread the data evenly. If the key is “city”, the Tehran Shard becomes much bigger than the others. This is called a Hot Shard.
- Common queries must go to only one Shard. “This customer’s orders” with the customer ID key reads only one Shard. But “all of today’s orders” must read from all Shards, and it is slow.
- We have no transactions across Shards. Things that must change together must be in one Shard.
- Adding a Shard is hard. With simple remainder division, when the number of Shards goes from three to four, most keys move. The Consistent Hashing method makes this movement much smaller.
Code
Atomic stock decrease in the database
The stock condition is inside the same UPDATE statement. The database runs this one statement atomically.
app.MapPost("/sale/{productId:guid}/buy", async (
Guid productId, ShopDbContext db, CancellationToken ct) =>
{
// One atomic statement: check and decrement together.
int updated = await db.Products
.Where(p => p.Id == productId && p.Stock > 0)
.ExecuteUpdateAsync(s => s.SetProperty(p => p.Stock, p => p.Stock - 1), ct);
return updated == 1
? Results.Accepted()
: Results.Conflict("Sold out");
});
The same thing with Redis
For very heavy load, we keep the stock in Redis. Redis runs commands one by one, so two DECR commands never run in the middle of each other.
long left = await redis.StringDecrementAsync($"sale:{productId}:stock");
if (left < 0)
{
return Results.Conflict("Sold out");
}
await queue.PublishAsync(new PhoneReserved(productId, userId), ct);
return Results.Accepted();
For the “one phone per user” condition, the user check and the stock decrease must happen together. For example, with a Lua script in Redis or a database transaction. Otherwise, the stock may go down, but the purchase is rejected because the user is a repeat buyer.
Choosing the Shard
using System.Buffers.Binary;
public static int ShardFor(Guid customerId, int shardCount)
{
// Do NOT use GetHashCode(): it is not guaranteed to be stable,
// and string hash codes change in every process. The shard must never change.
Span<byte> bytes = stackalloc byte[16];
customerId.TryWriteBytes(bytes);
uint hash = BinaryPrimitives.ReadUInt32LittleEndian(bytes);
return (int)(hash % (uint)shardCount);
}
In .NET, the hash value of strings is different in each run of the program. If you choose the Shard with it, after a restart you look for the customer’s data in the wrong Shard.
Important rules
- Ask first, then design. How many users? How many requests per second? Is payment also in this flow? The numbers change the design.
- A stateless app. State goes in Redis or the database, not in server memory.
- Reduce the load early. CDN, cache and request limits before the server and the database.
- Sensitive work is atomic. Check and change in one action.
- Heavy work goes through a queue. The user does not wait for the invoice and the email to be made.
- Run more than one of each part. A single Load Balancer, a single Redis or a single database is a single point of failure.
- Scale the database in order. First indexes and cache, then Replicas, and last, Sharding.
Common mistakes
| Mistake | Result | The right way |
|---|---|---|
| Reading the stock and decreasing it in two separate statements | Two people buy the last phone. | One atomic statement, in the database or Redis. |
| Keeping the session in server memory | With horizontal Scale, the user loses their data. | State in Redis or the database. |
| The only answer: “more servers” | The database is the bottleneck, and more servers add more load to it. | Reduce the load before the database, and use a queue. |
| Reading from a Replica right after a write | The user does not see their own change. | Read your own writes from the Primary. |
| An uneven Shard key, like city | One Shard stays under load, and the rest are idle. | A key with even spread, like the customer ID. |
| Choosing the Shard with the GetHashCode method | After a restart, the data is searched for in the wrong Shard. | A fixed hash that does not depend on the run. |
Which tool for which problem?
These first
- Many reads: cache and Replicas.
- Sudden load: CDN, request limits and a queue.
- API servers under load: horizontal Scale with a stateless app.
Only with a reason
- Splitting data with Sharding: only when the data size or the writes are really more than one server can handle.
- The Hash method in the Load Balancer: only when removing state from the server is not possible.
Summary in seven lines
- Before you design, calculate the numbers: users, requests per second, and which part really does the work.
- A bigger server has a limit. More servers have much less of a limit, but the app must be stateless.
- A Load Balancer spreads the load and takes an unhealthy server out with a health check.
- Reduce the load early and give heavy work to a queue.
- Sensitive work, like decreasing the stock, must be atomic.
- With Replicas, reads are spread. Watch out for the copy delay.
- With Sharding, data is split. The right key is the most important decision, and this is the last option.