Core Building Blocks

Learn the essential infrastructure components used repeatedly in system design solutions.

Last generated

Lesson 4 of 18 available17 practice questions

SPACED REPETITION · 17 practice questions

Make this lesson stick.

Try 3 questions now. No account needed. Sample answers aren't saved.

One box, launch day

You are building a secondhand marketplace. People photograph things they want to sell, post listings, and browse everyone else's. Version one runs on one machine: the web app, PostgreSQL and the uploaded photos all live on a 16-core VM with a 2 TB SSD and a 10 Gbit/s network card. The beta runs fine on it.

Then marketing does its job. Here are the assumptions for launch month and what they imply:

Assumption Arithmetic Result
4 million daily active users, 50 API calls each 200 M calls ÷ 86,400 s 2,315 requests/s on average
The evening peak is 3× the average 2,315 × 3 about 6,900 requests/s
100,000 new listings a day, 6 photos each, 3 MB per original 100,000 × 6 × 3 MB 1.8 TB/day of originals
Three resized copies per photo, 400 KB in total 600,000 × 400 KB 240 GB/day more
Each user loads 200 images a day at 150 KB 4 M × 200 × 150 KB = 120 TB/day 11.1 Gbit/s on average, 33 Gbit/s at peak
New listings written to the database 100,000 ÷ 86,400 s about 1.2 inserts/s
Each API call needs 8 ms of CPU on the 16 cores 16 × 1,000 ms ÷ 8 ms at most 2,000 requests/s
One PostgreSQL machine serves simple indexed reads assumed for this design about 10,000 reads/s

Predict first: which parts of the one box fail on launch day, in what order, and which part that people usually blame first survives?

Check your answer

The app's CPU and the network card go first, in the first busy hour. The CPU tops out at 2,000 requests/s, which is below even the 2,315/s average, and the evening peak asks for 3.5 times that. The images need 11.1 Gbit/s on an average day, more than the 10 Gbit/s the card can carry. The disk follows before the day is out: originals plus resized copies add about 2 TB a day, so the 2 TB SSD is full in about 23.5 hours.

The database, the usual first suspect, survives. It absorbs 1.2 listing inserts a second, and even if every API call read from it, 6,900 simple indexed reads a second is within the assumed 10,000.

In a product full of photos, the app servers and the bytes (network and disk) break the design long before the queries do. And one problem shows up in no traffic number at all: the machine is a single point of failure. A kernel update that needs a reboot takes the whole marketplace offline.

Each of these failures has a standard fix, and each fix is a building block: a component whose behavior, limits and failure modes are well known, so you can use it without inventing it.

Pressure in the requirements Building block What you give up Taught in
More requests than one server can serve; any server can die Load balancer + stateless service tier An extra network hop; the balancer tier must itself be redundant this lesson
Large files: photos, video, documents, backups Object storage, usually behind a CDN Whole-object writes only; slow first byte without a CDN this lesson
Many servers must create records with unique keys, without one central counter ID generation scheme Clock dependence, 64 vs 128 bits, IDs that leak creation time this lesson
The same data is read over and over Cache Stale reads, invalidation work, memory Databases & Storage
Data must survive and be queried Database The hardest tier to scale; replicas lag Databases & Storage
Slow or spiky work the user need not wait for Queue + workers Work finishes later; duplicates and ordering to handle Messaging & Queues

This lesson teaches the first three in depth and ends with the map of where the other three sit. It assumes the life of a request and the estimation habits from Foundations of System Design, and the case for scaling out from Key Concepts & Terminology.

💡 Habit worth building: name the pressure before you name the box. "Photos arrive at 1.8 TB a day, so the bytes go to object storage" earns more credit than "I'll add S3". In an interview, every block you draw should come with three things: the pressure that put it there, the move it makes against that pressure, and what it costs.

Load balancers: one address, many servers

Pressure: more requests than one server can handle, and servers that die.

A load balancer gives clients one stable address and spreads their requests across a pool of interchangeable servers (the backends, or targets). It has three jobs:

  1. Spread the load so no backend drowns while another idles.
  2. Detect dead backends and send traffic elsewhere.
  3. Decouple clients from the pool, so you can add, remove and replace servers without anyone noticing.
                  clients
                     |
            api.market.example
                     |
             +---------------+
             | load balancer |
             +---------------+
             /       |       \
      +-------+  +-------+  +-------+
      | app 1 |  | app 2 |  | app 3 |
      +-------+  +-------+  +-------+

L4 or L7: how much of the request can it see?

Balancers differ in how much of the traffic they read.

  • An L4 (transport-layer) balancer sees TCP or UDP: addresses and ports. It chooses a backend when a connection starts and sends every packet of that connection to the same backend. It never parses the HTTP inside. Packet-forwarding designs such as Linux IPVS and Google's Maglev do not even end the client's TCP connection; the backend does. AWS Network Load Balancer works the same way unless you ask it to terminate TLS.
  • An L7 (application-layer) balancer ends the client's connection, reads each HTTP request, and sends it over its own connection to a backend. Because it sees the method, path, headers and cookies, it can send /api/* and /images/* to different pools, choose a backend per request, and retry a failed request elsewhere. AWS Application Load Balancer, NGINX, Envoy and HAProxy in HTTP mode work this way.

Follow one HTTPS request through each:

Step L4, packet-forwarding L7 proxy
Client opens TCP and TLS Backend chosen from the connection's addresses and ports; TLS usually passes through to it The balancer completes TCP and TLS itself
Client sends GET /listings/42 Forwarded as opaque bytes Decrypted and parsed; a backend is chosen for this request
Next request on the same connection Same backend, no new choice A new choice
The backend answers 503 Passed straight to the client Can retry on another backend, if the request is safe to repeat

What each costs you:

  • L4 is cheap per packet, handles enormous throughput, and works for any TCP or UDP protocol (databases, MQTT, game traffic). But it balances connections, not requests, and it cannot route by URL or retry.
  • L7 gives you content-based routing and per-request balancing. It pays with CPU for parsing and TLS on every request, with a proxy hop (two TCP connections instead of one), and it becomes a place that holds your TLS keys and sees plaintext.

Predict: four API-gateway instances call an internal recommendation service over gRPC. The service runs 20 pods behind an L4 balancer. gRPC uses HTTP/2, which keeps one long-lived connection per client and sends every request over it. What does the CPU graph of the 20 pods look like?

Check your answer

At most 4 pods are busy, one per gateway connection, and 16 sit idle. If two connections landed on the same pod, only 3 are busy. The L4 balancer did its job, because it balanced connections, but there are only four connections to balance.

The fix is balancing per request: an L7 proxy that understands HTTP/2 (Envoy, for example), or client-side balancing in the gRPC library using the list of pods from service discovery. The Kubernetes blog has described exactly this failure: a default Kubernetes Service balances at L4, so gRPC traffic piles onto a single pod (Kubernetes blog, 2018).

The same trap catches every long-lived connection. WebSockets stay on the server they first reached, so after a scale-out the new servers receive only new connections while the old ones stay hot.

Reverse proxy, L7 balancer, API gateway: one family

These terms blur because they describe the same kind of machine configured for different jobs.

  • A reverse proxy accepts connections on behalf of servers and forwards requests to them. Clients see the proxy, never the backends. Typical extras: TLS termination, compression, response caching, and buffering slow clients so a backend worker is not tied up by someone on a bad mobile connection.
  • An L7 load balancer is a reverse proxy whose main job is spreading requests over a pool and routing around dead members.
  • An API gateway is a reverse proxy at the edge of a set of services, with API concerns added: authentication, per-client quotas and rate limits (rate limiters get their own lesson, Design URL Shortener & Rate Limiter), routing /orders and /users to different services, and sometimes combining several backend calls into one response.
  • An L4 packet-forwarding balancer is not a proxy, because it never ends the client's connection.
  • A forward proxy is the mirror image: it acts for clients going out, such as a corporate egress proxy.

NGINX, Envoy and HAProxy can each be configured as a reverse proxy, an L7 balancer or a gateway, and each can also run as an L4 proxy that ends the TCP connection; the role comes from the configuration, not the product. Packet-forwarding L4 is the territory of IPVS, Maglev and NLB. Terminating TLS at the balancer is normal practice, and both ALB and NLB support it. If policy requires encryption inside your network, the balancer re-encrypts to the backends. If the backends themselves must hold the keys, use L4 pass-through.

A common layout: DNS sends image and static-file requests to a CDN, and API calls to a load balancer → API gateway → services. You will also see a gateway fleet acting as the edge proxy itself, with no separate balancer in front. Both are sound. The balancer spreads load over identical instances; the gateway decides which service a request belongs to.

⚠️ Keep the gateway thin. It sits on every request's path, so business logic there turns it into a bottleneck and a shared deploy risk. It must scale out and be redundant like any other tier.

Choosing a backend

Algorithm How it picks Good when Weakness
Round robin Next in turn Backends are equal and requests cost about the same Blind to how busy each backend is
Weighted round robin In turn, in proportion to weights Mixed machine sizes; sending 5% to a canary Weights are static
Random Uniformly at random Many balancer nodes with no shared state Uneven over short windows
Least connections / least outstanding requests Fewest requests in flight Request costs vary widely Needs fresh counts; each balancer node sees only its own traffic
Power of two choices Samples two backends at random, takes the less busy Many balancer nodes, stale counts Worse than exact least-outstanding with perfect information (see the simulation), far better than round robin, and robust to stale counts
Hash of a key (IP, cookie, user ID) The same key always goes to the same backend Affinity: warm per-user caches, per-key state Skewed keys make hot backends; resizing remaps keys

Real products offer these directly. AWS Application Load Balancer offers round robin (the default), least outstanding requests and weighted random (ALB docs); Envoy's least-request balancer uses power of two choices by default.

Does the choice matter? Simulate it. Ten servers each handle one request at a time. Nine requests in ten are cheap (5 ms, a listing fetch) and one in ten is expensive (100 ms, a search). Requests arrive at random, keeping the servers 80% busy on average. Every policy sees exactly the same stream.

import heapq
import random

def simulate(policy, servers=10, requests=200_000, load=0.8, seed=7):
    rng = random.Random(seed)
    cheap, costly = 5, 100                      # ms of work per request
    mean = 0.9 * cheap + 0.1 * costly           # 14.5 ms
    rate = load * servers / mean                # arrivals per ms
    arrivals, work, t = [], [], 0.0
    for _ in range(requests):                   # same workload for every policy
        t += rng.expovariate(rate)
        arrivals.append(t)
        work.append(cheap if rng.random() < 0.9 else costly)

    pick = random.Random(seed + 1)
    free_at = [0.0] * servers                   # when each server's queue drains
    busy = [[] for _ in range(servers)]         # finish times of unfinished requests
    latencies, turn = [], 0
    for now, w in zip(arrivals, work):
        for q in busy:                          # forget requests that have finished
            while q and q[0] <= now:
                heapq.heappop(q)
        if policy == "round robin":
            s = turn % servers
            turn += 1
        elif policy == "random":
            s = pick.randrange(servers)
        elif policy == "least outstanding":
            s = min(range(servers), key=lambda i: len(busy[i]))
        else:                                   # power of two choices
            a, b = pick.sample(range(servers), 2)
            s = a if len(busy[a]) <= len(busy[b]) else b
        done = max(now, free_at[s]) + w         # one request at a time, in order
        free_at[s] = done
        heapq.heappush(busy[s], done)
        latencies.append(done - now)
    latencies.sort()
    p = lambda q: latencies[int(q * len(latencies)) - 1]
    return p(0.50), p(0.99)

for policy in ["round robin", "random", "power of two choices", "least outstanding"]:
    p50, p99 = simulate(policy)
    print(f"{policy:22} p50 {p50:6.1f} ms   p99 {p99:6.1f} ms")
Policy p50 p99
Round robin 85.8 ms 575.8 ms
Random 100.0 ms 730.0 ms
Power of two choices 13.4 ms 273.9 ms
Least outstanding 5.0 ms 191.3 ms

Other seeds move individual numbers by up to about 25% (round robin's p99 ranged from 576 to 728 ms over 30 seeds), but the order never changes. Real servers run several requests at once, which narrows the gaps: rerun with four request slots per server and p99 is 167 ms for round robin against 100 ms for least outstanding; with sixteen slots the gap almost disappears. The algorithm matters most when each backend can run only a few requests at a time (single-threaded runtimes, GPU inference, small connection pools) and when request costs are very uneven.

Why it works. A request's latency is mostly time spent waiting behind other requests. Round robin hands the next request to a server stuck behind a 100 ms search while other servers sit idle. Least outstanding sends it to an idle server whenever one exists, so the median request does not wait at all. Two random samples are enough to avoid the longest queues most of the time, and that matters in production, where you run many balancer nodes and each sees only its own in-flight counts. If they all chase "the least loaded backend" using stale counts, they pile onto the same one together; two random choices spread them out.

🧠 Say it like this: "Request costs vary a lot here, searches versus point reads, so I'd use least outstanding requests, or power of two choices across a fleet of balancers. Round robin would queue cheap requests behind expensive ones." A likely follow-up: "What if one backend is a bigger machine?" Weights, or least outstanding, which adapts on its own because the bigger machine finishes requests sooner.

Affinity: when the same key must reach the same server

Sometimes you want a key to land on the same backend every time: a server holding a large per-user cache, or a shard of in-memory state. Hash the key and pick a server from the hash.

The naive hash(key) % N has a nasty property. When N changes, almost every key moves: going from N to N + 1 servers remaps N/(N + 1) of the keys. From 9 servers to 10 that is 90% of the keys, and every moved key starts cold on its new server. Consistent hashing moves only the new server's share, about 1/10 here (9.6% in a 200,000-key run). How the ring and its virtual nodes work belongs to Advanced Topics & Final Prep. Also watch the key you hash: hashing by source IP puts an entire office, or a mobile carrier's shared address, onto one backend.

Health checks: noticing that a server died

Balancers learn about dead backends in two ways:

  • Active checks. The balancer probes each backend on a schedule, such as GET /healthz every 10 seconds with a 2-second timeout. After a number of consecutive failures (the unhealthy threshold) it stops sending traffic; after a number of consecutive successes (the healthy threshold) it resumes. Requiring several consecutive results both to leave and to return gives hysteresis, so one lost probe doesn't bounce a server in and out of the pool. A threshold of 1 does exactly that bouncing, called flapping.
  • Passive checks, also called outlier detection. The balancer watches real traffic and ejects a backend after, say, five consecutive errors or connection failures. Open-source NGINX's max_fails and fail_timeout settings work this way. Passive checks react faster to hard failures, but only for backends that are receiving traffic.

Follow a hang with active checks every 10 s, a 2 s timeout and an unhealthy threshold of 3. The server freezes at t = 0, just after a successful probe:

Time Event Failures counted
0 s Probe succeeds; the server freezes right after 0
10–12 s Probe times out 1
20–22 s Probe times out 2
30–32 s Probe times out; server removed 3

For 32 seconds the balancer keeps sending this server its share of traffic. With the marketplace's peak of 6,944 requests/s spread over 9 servers (the fleet sized in "Sizing the tier" below), that is about 772 requests/s, so roughly 24,700 requests hang until their clients time out. Defaults can be slower: an AWS Application Load Balancer checks every 30 s with a 5 s timeout and removes a target after 2 failures (ALB health check docs), up to about 65 s. Shorten the window with faster probes, passive checks, and a per-request timeout at the balancer that retries idempotent requests on another backend.

Predict: to be thorough, your /healthz runs SELECT 1 against the database, with the probe settings above. The database fails over and is unreachable for two minutes, while most pages could still be served from the cache. What happens to the site?

Check your answer

Every backend fails its health check at the same moment, so about 30 seconds in the balancer removes all of them and the whole site answers 503, including pages that come from the cache and would have worked. A partial outage became a total one, and it outlasts the database outage, because each backend must then pass several checks in a row before it returns.

AWS's Application Load Balancer happens to fail open: when every target is unhealthy, it routes to all of them anyway. Many balancers don't, and a design should not depend on that safety net. The rule: a health check answers "can this process serve requests?" (it is running and warmed up). How the process's dependencies are doing belongs in monitoring and in graceful degradation, not in the balancer's routing decision. Overload doesn't belong in the check either: at a fleet-wide peak every backend would fail at once, so a busy server should shed excess requests itself (429 or 503) and stay in the pool. Kubernetes asks two different questions, liveness (is the process stuck and in need of a restart?) and readiness (should it receive traffic right now?); keep dependency checks out of both.

Two more failure modes to name in an interview:

  • Draining. When you remove a backend on purpose (a deploy or a scale-in), the balancer should stop sending it new requests and let in-flight ones finish. AWS calls this the deregistration delay; it defaults to 300 seconds. Long-lived connections make draining slow, so WebSocket servers close connections in batches and let clients reconnect elsewhere.
  • Retry storms. Retries multiply across layers. If the client, the balancer and the service each try three times, one slow database query can turn into 3 × 3 × 3 = 27 attempts at the bottom, exactly when the database is least able to cope. Retry at one layer, cap retries with a budget, back off with jitter, and never blindly retry a request that is not idempotent, such as a payment (Networking Basics covers idempotency).

Who balances the balancer?

A single balancer is a single point of failure in front of everything you made redundant. The standard answers:

  • Active-passive pair with a floating IP. Two balancer machines run VRRP (keepalived is the usual tool). The active one owns a virtual IP address. If its heartbeats stop (one per second by default, with takeover after about three missed), the standby claims the address. Connections that were open through the failed machine are normally lost, and clients reconnect.
  • Active-active behind DNS or anycast. Several balancers each carry traffic. DNS returns several addresses, or one address is announced from several places (anycast) and routers spread the flows. Capacity grows with each balancer you add. DNS-based failover is only as fast as clients honor the record's TTL (see Networking Basics).
  • Managed cloud balancers run as a fleet across availability zones and scale themselves. Your job is to spread your backends across zones too.
  • Across regions, DNS that answers by geography or measured latency, or anycast, sends each user to a nearby healthy region, and each region runs its own balancers.
                    clients
                       |
        virtual IP 203.0.113.10 (VRRP)
                       |
          +------------+------------+
          |                         |
   +-------------+          +---------------+
   | LB-a ACTIVE |          | LB-b STANDBY  |
   | owns the IP |<-------->| takes the IP  |
   +-------------+ heartbeat| if LB-a dies  |
          |                 +---------------+
          v
   app servers (both LBs know the whole pool)

Your turn: sellers can now chat live with buyers over WebSockets: 300,000 concurrent connections at peak across 20 chat servers. Which balancer and algorithm do you choose, and what happens during a deploy and after a scale-out?

Check your answer

Either an L4 balancer or an L7 balancer that supports the WebSocket upgrade works; L7 lets you route /chat separately from the API on the same domain. Use least connections: each connection is long-lived, so the number of open connections is the load, and round robin over new connections ignores who is already full.

Deploys: draining could take hours if you wait for users to leave, so each server closes its connections in batches and clients reconnect, with jittered backoff so 15,000 of them don't return in the same second. Scale-out: new servers get only new connections, so the old ones stay hot until connections churn; ask some clients to reconnect to rebalance. How a message finds the recipient's server is chat design, covered in Chat Application (WhatsApp).

The stateless service tier

Pressure: you must add and remove servers freely, for peaks, deploys and failures, and the balancer must be able to send any request to any server.

A service tier is stateless when no request depends on something that only one instance holds: no session in local memory, no uploaded file on local disk, no in-memory counter that is the only copy. Instances still hold state, such as caches and connection pools, but only copies that can be rebuilt. That is what lets the balancer treat servers as interchangeable, and it is what makes autoscaling and rolling deploys routine.

Predict: on the one box, login stored the session in the app's memory. You add a second server behind a round-robin balancer. What do users see?

Check your answer

They log in on server 1. Their next request lands on server 2, which has never heard of them, so they appear logged out on every other click. The same bug wears other costumes: a photo saved to server 1's disk returns 404 when server 2 is asked for it, and a cron job defined on every server sends every reminder email twice.

Where the state goes

State On one box In a stateless tier What it costs
Login session App memory A shared session store (Redis with a replica, say) keyed by a session cookie, or a signed token One network lookup per request, usually under a millisecond in the same zone; or tokens that are hard to revoke
Uploaded files Local disk Object storage (next section) About 100–200 ms to the first byte; put a CDN in front for reads
Carts and drafts App memory The database or the session store A write per change
Slow background work A thread in the app A queue and workers The work finishes later
Scheduled jobs The box's crontab A scheduler that runs each job once, not once per instance One more component
Hot lookups In-process cache Still fine, as a copy with a TTL Each instance warms up separately and can briefly disagree with the others

Signed tokens (a JWT in a cookie or header) let the server keep nothing: it checks the signature and trusts the claims inside. The price is revocation. A stolen token, or one belonging to a user who just logged out, stays valid until it expires. The usual compromise is short-lived access tokens (minutes) plus a refresh token that is checked against a store.

Sticky sessions: the tempting shortcut

A balancer can pin each user to one server with a cookie (AWS calls this stickiness). It hides the state problem without solving it:

  • when the server dies, its users lose their sessions;
  • deploys and scale-ins must either wait for sessions to end or break them;
  • load becomes uneven: one server collects a few very heavy users, and after a scale-out the existing users stay on the old servers.

Stickiness is fine as a performance hint when the state can be rebuilt, such as a warm per-user cache. It is a trap when the pinned server holds the only copy.

Sizing the tier

How many app servers does the marketplace need at peak? Assumptions: 6,944 requests/s at peak; 8 ms of CPU per request; 16 cores per server; plan for at most 60% CPU so a spike or a lost server doesn't tip you over; three availability zones, and you must survive losing one.

Step Arithmetic Result
Ceiling per server 16 cores × 1,000 ms ÷ 8 ms 2,000 requests/s
Planned load per server 2,000 × 0.6 1,200 requests/s
Servers for the peak 6,944 ÷ 1,200 = 5.8 6
Survive losing one of three zones the remaining two zones must hold 6, so 3 per zone 9 servers
Requests in flight per server 1,200/s × 60 ms average response time about 72

The last row is Little's law: requests in flight = arrival rate × time in the system. A request that spends 60 ms mostly waiting on the cache and the database still occupies a worker for those 60 ms. A pool of 32 threads therefore caps a server at 32 ÷ 0.06 s ≈ 533 requests/s while its CPU sits mostly idle. Size worker pools, or use asynchronous I/O, from this number rather than from the core count.

⚠️ Autoscaling is not instant. A new instance takes minutes to boot, pass health checks and warm its caches. That is why you keep headroom (the 60%) and scale on a leading signal such as requests per instance, not only on CPU. Some balancers also ramp a new backend up gradually instead of handing it a full share of traffic cold.

Your turn: product wants "log out of all devices" to take effect within 5 seconds, and the platform team refuses to add a network lookup to every request. Sessions in Redis add a lookup; signed tokens that live 15 minutes break the 5-second rule. What do you propose?

Check your answer

Keep signed, short-lived tokens, and add a revocation list held in memory on every instance. When a user logs out everywhere, delete the user's refresh tokens in the store, so no new access tokens can be minted, then write "access tokens for user 81 issued before 14:03:07 are invalid" to a small store and broadcast it through pub/sub. Each instance checks incoming tokens against its in-memory copy, so no request makes a network call, and revocation spreads in about a second.

What you give up: every instance holds the list, which stays small because entries can be dropped once the tokens they cover have expired. An instance that misses a broadcast must catch up, so instances also re-read the full list periodically and at startup. A session store with a lookup per request is the simpler answer if the platform team relents. Say that too.

Object storage: where the bytes live

Pressure: large blobs (photos, video, PDFs, backups, logs) that grow without limit, are written once and read many times.

Object storage (Amazon S3, Google Cloud Storage, Azure Blob Storage, and S3-compatible systems such as MinIO and Ceph) stores objects: a blob of bytes plus metadata, addressed by a key inside a bucket, and read and written over HTTP.

What you get (S3's documented behavior):

  • Capacity and object counts that are effectively unlimited. You pay per GB stored per month, per request, and per GB sent out.
  • Durability designed for 99.999999999% ("eleven nines"), achieved by storing the data redundantly across several facilities.
  • Strong read-after-write consistency: once a PUT succeeds, every later GET or LIST sees it (S3 consistency model).
  • Request rates that scale by key prefix: at least 3,500 writes and 5,500 reads per second per prefix, with S3 adding capacity as load grows.

What you give up:

  • It is not a filesystem. You write an object whole; you don't edit bytes in the middle of it. A "rename" is a copy plus a delete, and "folders" are only key prefixes.
  • Latency. AWS quotes roughly 100–200 ms for small objects or the first byte of a large one. That is fine for images behind a CDN and wrong for a hot key-value lookup.
  • No queries. You can list keys by prefix and that is all. Anything you filter or sort by lives in a database.
Block storage (a cloud disk) File storage (NFS) Object storage
Accessed by Usually one server that mounts it Many servers sharing a filesystem Any client, over HTTP
Change one byte Yes Yes No, write a new object
Grows to One volume's size One shared filesystem Effectively unlimited
Typical use Database data files Shared files for existing apps Media, backups, logs, data lakes

The pattern: facts in the database, bytes in the bucket

The database row holds what you query: photo ID, listing ID, owner, status, and the object key, such as listings/42/7156281536.jpg. The bucket holds the bytes. Putting 3 MB blobs in the database bloats its backups and replication; keeping queryable facts only inside objects makes them unsearchable.

Store the key, not the URL. Build URLs when you read, from configuration: the CDN's domain plus the key. When you change CDN provider, rename the domain or move buckets, you edit one setting instead of rewriting 600 million rows.

The upload path: let clients talk to the bucket

client           API         database     object store   resize worker
  |               |              |              |              |
  | 1 add photo   |              |              |              |
  |-------------->|              |              |              |
  |               | row: pending |              |              |
  |               |------------->|              |              |
  | 2 presigned   |              |              |              |
  |   PUT URL     |              |              |              |
  |<--------------|              |              |              |
  | 3 PUT 3 MB    |              |              |              |
  |-------------------------------------------->|              |
  |               |              |              | 4 created    |
  |               |              |              |   event,     |
  |               |              |              |   via queue  |
  |               |              |              |------------->|
  |               |              |              | 5 resized    |
  |               |              |              |   copies     |
  |               |              |              |<-------------|
  |               |              | 6 ready      |              |
  |               |              |<----------------------------|
  1. The client asks to add a photo. The API checks that the user may edit listing 42 and inserts a photo row with status pending and a key.
  2. The API returns a presigned URL: a URL signed with the server's credentials that allows one kind of operation (a PUT to this key) until it expires, here in 15 minutes.
  3. The client uploads straight to object storage. The 3 MB never pass through your app servers.
  4. Object storage emits an "object created" event, and a queue delivers it to a resize worker.
  5. The worker writes the thumbnail, card and full-size copies back to object storage, under keys derived from the photo ID.
  6. The worker marks the row ready. Readers now receive CDN URLs built from the keys.

Why bother? Routed through the app tier, uploads would carry 1.8 TB a day, about 170 Mbit/s on average and 500 Mbit/s at peak, through servers sized for small JSON requests. Worse, each slow phone upload holds a worker. At the peak of about 21 photos per second, 20-second uploads keep roughly 420 requests in flight, about 46 per server, most of the 72-request budget computed above.

Failure modes, and what each costs:

  • The client never uploads. The row stays pending forever, so a cleanup job deletes rows still pending after a day.
  • The photo lands but the worker never finishes. The event is lost or the worker keeps crashing, so the bytes sit in the bucket while the row says pending. Queue retries and a dead-letter queue handle most of it; the same cleanup job deletes the object together with the stale row, so nothing is orphaned.
  • A presigned URL leaks. It is a bearer token: whoever holds it can use it, even repeatedly, until it expires. Keep expiry short, sign one key per URL, and check the object's type and size before marking it ready. S3's presigned POST policies can also cap the size with a content-length-range condition.
  • A huge file. One S3 PUT accepts up to 5 GB; AWS recommends multipart upload from about 100 MB. Multipart splits the file into parts of 5 MiB to 5 GiB, up to 10,000 of them, uploaded in parallel and retried one by one, for objects up to about 50 TB (S3 upload limits). A video on a flaky connection resumes from the last good part instead of starting over.

Serving: put a CDN in front

At peak the marketplace serves about 27,800 image requests a second, 33 Gbit/s. You don't serve that from the bucket. A CDN caches images near users; at a 95% hit ratio the bucket sees about 1,400 requests/s and 1.7 Gbit/s at peak. Because a replaced photo gets a new key, an image's bytes never change under its URL, so the CDN can cache it for a very long time and never needs invalidating for updates. Deletions are different: a takedown or a seller removing a photo of their home still needs a CDN purge, or a TTL short enough to live with. Private files, such as invoices, use short-lived signed URLs instead of public caching. Pull versus push CDNs, cache headers and invalidation are in Networking Basics.

Your turn: sellers may now attach a 2–10 minute video (up to 2 GB) to a listing. The app servers keep a 30-second request timeout. Sketch the upload and what happens after it.

Check your answer

The API creates a video row with status pending and starts a multipart upload, handing the client presigned URLs for the parts. The client uploads parts in parallel straight to object storage, retries only the parts that fail, and asks the API to complete the upload when all parts are in. None of that touches the 30-second timeout, because the app servers only hand out URLs.

The "object created" event queues a transcoding job; the row moves from pending to processing to ready, and the listing page shows a placeholder until then. A lifecycle rule aborts multipart uploads left incomplete for a few days, so abandoned parts don't pile up as storage you pay for. Playback goes through the CDN; adaptive bitrate streaming is covered in YouTube / Netflix Streaming.

Unique IDs without a single counter

Pressure: many servers create records at the same time, and each record needs a unique key without a round trip to one central counter.

On the one box, PostgreSQL's bigserial hands out 1, 2, 3, and that goes surprisingly far: one primary database can issue thousands of IDs a second. It stops being enough when:

  • writes are split across databases, and two shards would both hand out 1,001;
  • you need the ID before the row exists: the photo key in the upload flow, a record created offline on a phone, a message ID the client retries with;
  • writers sit in several regions, and a cross-region call on every insert costs tens of milliseconds.

Decide what the ID must do before choosing how to make it:

Property Why it matters
Unique The whole point
64 or 128 bits 64 fits a BIGINT and halves the key's size in every index; JavaScript numbers are exact only up to 2⁵³, so send 64-bit IDs to browsers as strings
Roughly time-ordered "Newest first" becomes an index scan, and B-tree inserts land at the right edge
No coordination per ID No network call per insert, and no outage when a coordinator is down
Reveals nothing Sequential IDs leak your volume and invite scraping /orders/1001, /orders/1002…

The options:

Scheme How it works Time-ordered? Coordination Watch out for
Database auto-increment 1, 2, 3… from one database Yes Every insert goes through it One writer; leaks counts
Offsets per shard Shard k of n issues k, k + n, k + 2n… Within a shard None after setup Adding shards means re-planning the scheme
Ticket server A database whose only job is issuing numbers Yes with one server; roughly with an odd/even pair One call per ID Must itself be redundant
Range leasing Each app server leases a block of 1,000 IDs at a time Roughly One call per block A crash throws away the rest of its block (gaps)
UUIDv4 122 random bits in 128 No None 16 bytes; inserts scatter across the index
UUIDv7 48-bit millisecond timestamp, then 74 bits of randomness Yes, to the millisecond None 16 bytes; reveals creation time
Snowflake-style 41-bit ms timestamp, 10-bit worker ID, 12-bit sequence, in 64 bits Roughly Worker IDs assigned once Clocks that go backwards; duplicate worker IDs

Flickr has described its ticket servers: two MySQL servers, one issuing odd numbers and one even, so that neither is a single point of failure (Flickr engineering, 2010). UUIDv7 was standardized in RFC 9562 in May 2024; PostgreSQL 18 ships uuidv7() and Python 3.14 ships uuid.uuid7().

Are random UUIDs safe from collisions? With n IDs drawn from 2¹²² values, the chance of any collision is about n² ÷ 2¹²³. For a trillion IDs that is about 9 × 10⁻¹⁴. A broken random number generator is a far likelier source of duplicates than the math.

Snowflake: 64 bits, three fields

Twitter has described Snowflake, its ID service from 2010: a 41-bit millisecond timestamp counted from a custom epoch (about 69 years of range), a 10-bit machine ID (1,024 machines) and a 12-bit sequence number (Snowflake README).

 bit 63   bits 62..22                    bits 21..12         bits 11..0
 +----+-------------------------------+-------------------+------------------+
 | 0  | ms since custom epoch (41)    | worker ID (10)    | sequence (12)    |
 +----+-------------------------------+-------------------+------------------+

Here is a generator. The clock is a parameter so the demo is deterministic:

import threading
import time

EPOCH_MS = 1_704_067_200_000        # custom epoch: 2024-01-01 00:00:00 UTC
WORKER_BITS, SEQ_BITS = 10, 12
MAX_SEQ = (1 << SEQ_BITS) - 1       # 4095

def wall_clock_ms():
    return time.time_ns() // 1_000_000

class Snowflake:
    def __init__(self, worker_id, clock=wall_clock_ms):
        if not 0 <= worker_id < (1 << WORKER_BITS):
            raise ValueError("worker id must fit in 10 bits")
        self.worker_id, self.clock = worker_id, clock
        self.last_ms, self.seq = -1, 0
        self.lock = threading.Lock()

    def next_id(self):
        with self.lock:
            now = self.clock()
            if now < self.last_ms:                 # NTP stepped the clock back
                raise RuntimeError(f"clock went back {self.last_ms - now} ms")
            if now == self.last_ms:
                self.seq = (self.seq + 1) & MAX_SEQ
                if self.seq == 0:                  # 4,096 IDs this ms: wait
                    while now <= self.last_ms:
                        now = self.clock()
            else:
                self.seq = 0
            self.last_ms = now
            return ((now - EPOCH_MS) << (WORKER_BITS + SEQ_BITS)
                    | self.worker_id << SEQ_BITS
                    | self.seq)

def decode(snowflake_id):
    return {"ms_since_epoch": snowflake_id >> (WORKER_BITS + SEQ_BITS),
            "worker": (snowflake_id >> SEQ_BITS) & ((1 << WORKER_BITS) - 1),
            "seq": snowflake_id & MAX_SEQ}

# Deterministic demo: three calls in one millisecond, then one in the next.
ticks = iter([1_750_000_000_000] * 3 + [1_750_000_000_001])
gen = Snowflake(worker_id=37, clock=lambda: next(ticks))
for _ in range(4):
    i = gen.next_id()
    print(i, decode(i))

Output:

192656126771351552 {'ms_since_epoch': 45932800000, 'worker': 37, 'seq': 0}
192656126771351553 {'ms_since_epoch': 45932800000, 'worker': 37, 'seq': 1}
192656126771351554 {'ms_since_epoch': 45932800000, 'worker': 37, 'seq': 2}
192656126775545856 {'ms_since_epoch': 45932800001, 'worker': 37, 'seq': 0}

Within one millisecond the sequence counts up. In the next millisecond it resets and the ID jumps by 2²² − 2, because the timestamp sits above the other 22 bits. Numeric order is time order.

Why it works. Two IDs from one worker never collide while the process runs: the pair (timestamp, sequence) never repeats, because the sequence counts within a millisecond and the generator refuses to go back in time. The guard lives in memory, so a process that restarts after the clock stepped back (a VM restored from a snapshot, say) must first wait until the clock passes the last timestamp it persisted, or take a fresh worker ID. IDs from different workers differ in the worker bits. Across workers the order is only rough: if two servers' clocks disagree by 3 ms, an ID minted later on one can be smaller than an ID minted earlier on the other. Snowflake's own documentation promises only rough ordering (IDs "k-sorted"), not a strict global sequence.

The capacity math: 4,096 IDs per millisecond is about 4.1 million per second per worker, and 2⁴¹ ms is about 69.7 years, so an epoch of 2024 runs out in September 2093.

Failure modes

  • The clock goes backwards. An NTP correction or a VM migration steps the clock back 6 ms. Without the guard, the generator would reuse timestamps it already used, restart the sequence at 0 and mint duplicates. So it refuses, or waits out small steps and alerts on large ones. Configuring time sync to slew (adjust gradually) rather than step makes this rare.
  • Two workers share an ID. Two instances with worker ID 37 mint identical IDs whenever they issue one in the same millisecond at the same sequence number. Worker IDs must be assigned, from a lease in a coordination service such as etcd or ZooKeeper, a database table, or a stable ordinal from the orchestrator. Picking one at random at boot is the birthday problem: 40 instances choosing among 1,024 IDs collide with probability 54%.
  • More than 4,096 IDs in one millisecond. The generator waits for the next millisecond; that caps one worker at about 4.1 million IDs per second.
  • A hot range partition. Time-ordered IDs are excellent keys inside one database's B-tree, and a poor key for splitting data into ranges across machines, because every new write lands in the newest range. Hash the ID to choose a shard instead (hot keys and partitioning are in Databases & Storage).
  • Information leaks. Time-ordered IDs reveal when a record was created, and sequential ones reveal how many exist: a competitor places two orders a day apart and subtracts. If that matters, expose a separate random public ID. An unguessable ID is not access control either way; the server must still check that the caller may read /orders/….

Random keys in a B-tree

Most databases keep the primary key in a B-tree, and MySQL's InnoDB stores the whole table in primary-key order. Time-ordered keys arrive in increasing order, so inserts go to the right edge: the same few pages stay in memory and fill up completely. Random UUIDv4 keys land anywhere, so each insert touches a random page that may not be in memory, and pages split and stay partly empty. Once the index is bigger than memory, the difference shows up as disk reads on every insert. In InnoDB every secondary index also stores the primary key, so a 16-byte key instead of an 8-byte one makes every index bigger. Indexes in depth are in Databases & Storage.

🧠 Say it like this: "I'd use Snowflake-style 64-bit IDs. There's no coordination per insert, they fit a BIGINT, and they sort by time, so 'newest first' is an index scan. The costs are a dependency on the clock, which the generator guards against, and assigning worker IDs, which I'd lease from etcd." A likely follow-up: "Why not UUIDv7?" It would also work. It needs no worker IDs at all, and costs 16 bytes instead of 8.

Your turn: order IDs appear in customer-facing URLs, and the business does not want competitors estimating order volume. Orders live in one MySQL (InnoDB) primary, with a peak of 3,000 inserts a second and about 1,000 a second on average. What do you use?

Check your answer

Two IDs. Internally, a compact time-ordered key: a BIGINT auto-increment is fine at 3,000 inserts a second on one primary, or Snowflake-style if you expect to shard. Publicly, a random identifier (a UUIDv4, or 96+ random bits in base-62) with its own unique index, used in URLs and never the internal key.

UUIDv7 as the single ID is also defensible: its random bits reveal essentially nothing about volume, only creation time, and it is time-ordered in the index at 16 bytes per key. Plain auto-increment in URLs fails the requirement: two orders placed exactly a day apart whose IDs differ by 86 million tell a competitor you took about 86 million orders that day.

The map: where caches, databases and queues fit

Put the pieces together for the marketplace. Uploads, not drawn here, go straight from the client to object storage with presigned URLs.

   DNS answers "where is img.market.example / api.market.example?"
   and then steps aside: it carries none of the traffic below.

              browser / mobile app (its own HTTP cache)
                |                                  |
          image requests                       API calls
                |                                  |
            CDN edge                       load balancer (L7)
                | on a miss                        |
                v                        stateless app servers
         object storage                   (in-process cache)
           |        ^                     |         |         |
           |        |                   cache   database    queue <---+
           |        |                  (Redis)  (primary +    |       |
           |        |                            replicas)    v       |
           |        +------------------------------------- workers    |
           |             resized photos             (resize photos,   |
           |                                         send emails,     |
           |                                         update search)   |
           +-------------------- "object created" events -------------+

The app servers look in the cache first and go to the database on a miss (the cache-aside pattern). Two sources feed the queue: object storage publishes an "object created" event when an upload lands, which triggers resizing, and the API publishes jobs such as emails and search-index updates. The workers write resized photos back to object storage.

Caches: pick the layer

Layer Holds Cost of a hit Shared by How it gets fresh
Browser or app Images, scripts, styles No network One user Cache headers; new URL when content changes
CDN Images, static files, some public API responses A nearby edge Everyone near that edge TTL, purges, versioned keys
In-process Small hot data: categories, feature flags Microseconds One app instance TTL, per instance
Distributed cache (Redis, Memcached) Hot rows, sessions, computed results A network round trip, about a millisecond or less All app servers TTL and delete-on-write
Database buffer pool Hot pages Inside the database The database Automatic

Cache as close to the user as the data's audience allows: public and unchanging data goes to the CDN; per-user data can live in the user's browser (private caching) or in the app tier, never in a shared CDN cache; fast-changing data gets a short TTL or is deleted on write.

Predict: the listing API serves 6,944 reads a second at peak through Redis. The hit ratio rises from 90% to 99%. By what factor does database read load fall?

Check your answer

Ten. The database sees the misses: 10% of 6,944 is about 694 reads a second, and 1% is about 69. A hit ratio that looks only "9 points better" is a 10× change for the database. That is why you estimate the miss rate, not the hit rate. How to get there, and keep the cache correct (cache-aside, write-through, write-back, eviction, stampedes), is in Databases & Storage.

Databases: the source of truth

The map has two sources of truth: the database for facts and the original photos in object storage for bytes. Everything else is derived from them (cache entries, CDN copies, resized photos, search indexes) and can be rebuilt, while nothing can rebuild either source. The database is also the hardest tier to scale out, because every row has a home: read replicas carry read-heavy load (and lag behind), sharding splits large or write-heavy data. In the marketplace, one primary with replicas is plenty, because 1.2 listing writes a second and small metadata rows are all it carries. The bytes went to object storage. Choosing and scaling databases is Databases & Storage.

Queues: take slow work off the request path

Ask one question: does the user need this result before the response? If not, queue it. In the marketplace that covers resizing photos (seconds of CPU each, triggered by object storage's upload event), sending "your item sold" emails and updating the search index (published by the API). The request that caused the work answers in milliseconds; workers catch up, and a spike becomes a backlog instead of timeouts.

What you pay: the work finishes later (the listing shows a "processing" placeholder), a message may be delivered more than once, so workers must be idempotent (write resized copies to keys derived from the photo ID, so a repeat overwrites the same objects), and the backlog needs monitoring. Delivery guarantees, ordering and dead-letter queues are in Messaging & Queues.

Every box has a price

Every hop is another network round trip (well under a millisecond inside one data center, but they add up), another thing to deploy, monitor and page someone about, and another way to fail. The marketplace earned each block with a number; a prototype with a hundred users earns almost none of them. Interviewers reward "I'd add a cache if reads outgrow one database" over boxes drawn by reflex, and they discount "because big company X uses it".

Final round: no labels

Real prompts don't say which block they want. Find the pressure first, then pick the move and name its price.

Challenge 1: invoices at month end

A billing service must generate 1.2 million invoice PDFs on the first of each month (2 seconds of CPU each), email each customer a link, and keep every invoice downloadable for 7 years. The API tier is six small servers with a 30-second timeout. Name the blocks and the numbers that justify them.

Check your answer
  • Queue and workers. 1.2 M × 2 s is about 667 CPU-hours, which no request can wait for. The API enqueues one job per invoice; a worker pool with 200 cores finishes in about 3.3 hours. Workers must be idempotent, keyed by invoice ID, so a redelivered job overwrites the same PDF instead of producing a second invoice.
  • Object storage for the PDFs, under keys derived from the invoice ID. A lifecycle rule moves invoices older than a few months to a cheaper class that still serves reads in milliseconds (S3 Standard-IA or Glacier Instant Retrieval), since they are rarely read but must be kept. An archive class that needs a restore of minutes to hours first would break the instant download below.
  • Short-lived presigned download URLs, not a public CDN, because invoices are private. The email links to your app, which checks the user and then redirects to a fresh presigned URL, so an old email still works years later.

Challenge 2: the internal tool

An expense tool for 200 employees, at most 5 requests a second during office hours. It must be deployable without downtime and survive the loss of one server. Draw the design and say what you would add later, and when.

Check your answer

Two small app instances behind a managed load balancer, plus a managed PostgreSQL database with automated backups and a standby. Two instances, not one, because the requirements say no downtime on deploys and survive one server; no bigger number is needed.

No cache (5 requests a second is nothing), no queue until something slow appears (then add one for, say, report exports), no CDN unless remote offices complain about static files, and no ID scheme beyond the database's own. Receipt photos go to object storage from day one, because that is cheaper and simpler than disk, not because of scale. Saying what would trigger each addition shows more judgment than drawing them all.

Challenge 3: a flash sale

A ticketing site sells 50,000 concert tickets at 10:00 sharp. 400,000 people have the page open, and within ten seconds of 10:00 they send 400,000 "buy" requests, about 40,000 a second. You scaled the load balancer and the app tier ahead of time: 40 servers, enough for 48,000 requests a second. Which block fails first, and what do you put in front of it?

Check your answer

Not the load balancer and not the app tier: both were scaled for a known event. The database fails first. 50,000 seats are 50,000 heavily contended rows, and 40,000 reservation attempts a second on them become lock waits, timeouts and retries that make the pile-up worse.

Put a queue in front of the reservation: each click gets a place in line and a "you're in line" page, and a controlled number of workers reserve seats at a rate the database can sustain. Serve the event page itself from the CDN so the stampede of page loads never reaches your servers. Turn off balancer retries for the reservation call, because it is not idempotent unless each attempt carries an idempotency key. The price: buyers wait in a visible line instead of racing, which is also fairer. Contention and reservation design continue in Databases & Storage and Messaging & Queues.

Cheat sheet: pressure → move → cost

When the requirements say… Reach for And say what it costs
More requests than one server; servers die Load balancer over a stateless tier A hop; the balancer tier needs its own redundancy
Request costs vary widely Least outstanding requests, or power of two choices Needs in-flight counts; small overhead
Long-lived connections (gRPC, WebSockets) Per-request L7 balancing, or least connections L7 CPU; slow draining and rebalancing
The same key must hit the same server Consistent hashing Skewed keys still make hot servers
A dead server must leave the pool fast Short active checks plus passive checks; retries for idempotent calls False alarms if thresholds are too tight
Users must stay logged in across servers Session store, or signed short-lived tokens A lookup per request, or slow revocation
Big files Object storage, presigned uploads, CDN for reads Whole-object writes; 100–200 ms first byte at the origin
Unique IDs from many writers Snowflake-style 64-bit, or UUIDv7 Clock and worker-ID care, or 16-byte keys
Repeated reads A cache at the right layer Staleness and invalidation
Slow work the user needn't wait for A queue and idempotent workers Later completion, duplicate deliveries

Before moving on, explain the marketplace aloud as you would at a whiteboard, without notes: follow one photo upload and one listing view end to end, and at every box say the pressure that put it there and what it costs. Then answer three follow-ups: one app server dies mid-request; a server's clock steps back 50 ms; the active load balancer's machine loses power. If each answer names the mechanism that saves you (or doesn't), you are ready.

Next: Databases & Storage goes deep on the database and cache boxes, and Messaging & Queues on the queue. To see ID generation put to work in a complete design, try URL Shortener (e.g. TinyURL).