Levelwise
English
Architecture and system design

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.

Not reviewedWritten with AI helpReading time: 18 minExample of a phone flash sale in an online shopC# and .NET 10 code

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:

  1. No more than one thousand phones are sold.
  2. The system does not go down.
  3. 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.

  1. Two hundred thousand users in sixty seconds means about 3300 users per second.
  2. If each user clicks or refreshes five times on average, that means about 16 thousand requests per second.
  3. 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

Vertical: bigger serverHorizontal: more servers2 CPU32 CPUHas a limit; a single point of failureThe code does not changeLoad Balancer2 CPU2 CPU2 CPUAlmost no limit; one broken server is OKThe app must be stateless
A bigger server has a limit. More servers have almost no limit, but the app must be stateless.
  • 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:

  1. The user’s first request goes to server 1, and the shopping cart is saved in the memory of server 1.
  2. The same user’s second request goes to server 2.
  3. Server 2 does not see the cart. The user thinks the cart is now empty.
  4. 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:

Spreading load over three API copies
Server 1Healthy0
Server 2Healthy0
Server 3Healthy0
Send requests, then break a server and send again.

    Designing the flash sale

    Now let us put everything together:

    Users200K peopleCDN + Rate LimitMost load stops hereLoad BalancerAPI 1API 2API 3StatelessRedis: DECR stockAtomic stock decreaseQueueOrder WorkerPrimaryRead ReplicasCopyHeavy work at its own speed
    Each layer takes part of the load. Only the needed work reaches the database.
    1. 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.
    2. Stateless API servers, many of them. Behind the Load Balancer, with horizontal Scale.
    3. 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.
    4. 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.
    5. Protect the database. Send the reads to the Replicas. Writes go only to the Primary.
    Why are two separate actions dangerous? Imagine only one phone is left. Two requests read the stock at the same moment, and both see the number 1. Both decrease the stock. Now the stock is minus one, and one extra phone is sold. This is called a Race Condition.

    Database: Replica

    Most systems read much more than they write. A customer looks at a product a hundred times and buys once. So:

    Orders APIPrimaryWrites onlyReplica 1Replica 2Reads are spread hereWriteReadDelayed copyThe customer placed an order, but a read from a copy may not show it yet
    Writes go only to the Primary. Reads are spread between the Replicas. The copy has a small delay.
    1. One Primary takes all the writes.
    2. Several Replicas have a copy of the data, and the reads are spread between them.
    3. 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).

    Remember: With Replicas, only reads scale. All writes still go to one server. If there are many writes, Replicas do not help.

    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.

    New ordercustomerId = 7f3a…hash(customerId) % 3Which database?Shard 0Shard 1Shard 2Result: 1All orders of one customer in one databaseEach database has only one third of the data and load
    The customer ID decides which database the data is in.

    The most important decision is choosing the Shard key:

    1. 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.
    2. 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.
    3. We have no transactions across Shards. Things that must change together must be in one Shard.
    4. 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.
    Splitting data is the last option: First try the right index, cache, Replicas and a stronger server. Sharding adds a lot of complexity, and going back from it is almost impossible.

    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

    1. Ask first, then design. How many users? How many requests per second? Is payment also in this flow? The numbers change the design.
    2. A stateless app. State goes in Redis or the database, not in server memory.
    3. Reduce the load early. CDN, cache and request limits before the server and the database.
    4. Sensitive work is atomic. Check and change in one action.
    5. Heavy work goes through a queue. The user does not wait for the invoice and the email to be made.
    6. Run more than one of each part. A single Load Balancer, a single Redis or a single database is a single point of failure.
    7. 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

    1. Before you design, calculate the numbers: users, requests per second, and which part really does the work.
    2. A bigger server has a limit. More servers have much less of a limit, but the app must be stateless.
    3. A Load Balancer spreads the load and takes an unhealthy server out with a health check.
    4. Reduce the load early and give heavy work to a queue.
    5. Sensitive work, like decreasing the stock, must be atomic.
    6. With Replicas, reads are spread. Watch out for the copy delay.
    7. With Sharding, data is split. The right key is the most important decision, and this is the last option.