What this lesson gives you
In How systems grow you watched Recipebox, our recipe-sharing app, add components one problem at a time. This lesson is the reference shelf for those components. For each building block (a reusable component that appears in almost every large system) you get the same five short answers:
- What it is, in one or two plain sentences.
- Problem it solves.
- When to use it.
- When not to use it.
- Common interview follow-up, with a short model answer.
Before the blocks, we cover the handful of trade-off pairs that every design discussion comes back to: performance vs scalability, latency vs throughput, availability vs consistency. After the blocks, you will find consistency patterns, availability patterns with worked "nines" math, and communication choices (TCP vs UDP, RPC vs REST). The lesson ends with a one-table cheat sheet you can revise from the night before an interview.
Interviewers rarely ask "define a load balancer". They ask "why did you put a cache there, and what happens when it is wrong?" The follow-up answers in this lesson are written for exactly that moment. When you need more depth on one block, the existing advanced lessons in this track are linked.
Three trade-off pairs to understand first
Performance vs scalability
Performance is how quickly the system serves a single request when it is not busy. Scalability is the system's ability to keep that performance when load grows, usually by adding resources.
A simple test: if the system is slow for one user, you have a performance problem (a slow query, an inefficient algorithm). If it is fast for one user but slow under heavy load, you have a scalability problem (a shared resource that cannot be multiplied, such as one database primary).
Recipebox example: rendering a recipe page takes 2 seconds even at 3 a.m. because of a missing index. That is performance. The same page takes 80 ms at 3 a.m. but 3 seconds at the 8 p.m. peak. That is scalability.
Latency vs throughput
Latency is the time one operation takes, for example 120 ms from request to response. Throughput is how many operations the system completes per unit of time, for example 5,000 requests per second.
They are related but not the same. A highway analogy helps: latency is how long one car takes to drive from A to B; throughput is how many cars pass per hour. Adding lanes raises throughput without making any car faster. Batching work often raises throughput but increases latency for each item, because items wait for the batch to fill. Good designs aim for the highest throughput that still meets a latency target, usually stated as a percentile: "99% of requests finish within 300 ms" (the p99 latency).
Availability vs consistency
Availability means every request gets a non-error response. Consistency (in the distributed-systems sense) means every read sees the most recent write.
When machines cannot talk to each other because of a network partition (a broken link between parts of the system), a distributed data store must choose: refuse some requests until it can be sure the data is current (consistent, less available), or answer with possibly stale data (available, less consistent). This choice is the CAP theorem, covered properly in fundamentals. In practice, you choose per feature: Recipebox's like count can be stale (choose availability), but a payment or a unique username claim cannot (choose consistency).
Interview tip
Do not say "I choose CP" or "I choose AP" for a whole system. Say which data needs which: "Usernames must be strongly consistent; the feed and like counts can be eventually consistent." That shows you understand the trade-off is per operation.
Getting a request in: DNS, CDN, load balancers and proxies
DNS
What it is. The Domain Name System turns a name like recipebox.example into an IP address. Your device asks a resolver (usually run by your ISP or a public provider), which asks the authoritative name servers for the domain and caches the answer for a time set by the record's TTL (time to live).
Problem it solves. People remember names, not numbers, and services can move to new addresses without telling every user.
When to use it. Always; every internet service relies on it. In design you also use DNS for simple traffic steering: weighted records (send 10% of users to a new cluster), latency-based routing (send users to the nearest region) and failover records (switch to a backup if health checks fail).
When not to use it. Do not rely on DNS for fast failover. Resolvers and clients cache answers, sometimes longer than the TTL, so a change can take minutes to reach everyone. For fast failover inside a region, use a load balancer.
Common follow-up. "Why might users still hit a dead server after you changed DNS?" Because resolvers, operating systems and browsers cached the old answer; lowering the TTL in advance reduces the delay but costs more DNS lookups. Record types and resolution steps are covered in DNS.
CDN (push vs pull)
What it is. A content delivery network is a set of cache servers ("edges") in many cities that serve content from near the user. On a miss, the edge fetches from your server, called the origin.
Problem it solves. Distance and origin load. Serving a 3 MB recipe photo from an edge 20 km away is faster than from a data centre on another continent, and the origin sees only the misses.
There are two ways content gets onto a CDN:
| Pull CDN | Push CDN | |
|---|---|---|
| How content arrives | Edge fetches from origin on the first request, then caches it until expiry | You upload content to the CDN in advance |
| Good for | Sites with lots of content and unpredictable popularity | Content that changes rarely and is known in advance (release files, video catalogues) |
| Downside | First user in each region pays a slow miss; expired content refetched | You manage uploads and storage; you pay to store content even if nobody reads it |
When to use it. Static assets (JavaScript, CSS, images, video) and public pages that many people read.
When not to use it. Personalised or private responses, unless you are very careful with cache keys, and data that must be fresh to the second.
Common follow-up. "How do you update a file cached for a year?" Give each version a new name containing a content hash (app.3f9a1c.js), so new pages point to the new file and the old one simply stops being requested. More in networking and edge delivery.
Load balancer (L4/L7, algorithms, health checks)
What it is. A component that receives incoming connections or requests and spreads them across a pool of backend servers.
Problem it solves. One server cannot handle all traffic and can fail. A load balancer lets you add servers freely and routes around dead ones.
Layer 4 vs layer 7. The layers come from the OSI model (OSI and TCP/IP models):
- An L4 load balancer works at the transport layer. It sees IP addresses and ports, picks a backend per TCP connection, and forwards bytes without reading them. It is fast and protocol-agnostic.
- An L7 load balancer works at the application layer. It reads the HTTP request, so it can route
/api/searchto search servers and/imageselsewhere, add headers, terminate TLS and retry failed requests. It does more work per request.
Common algorithms:
| Algorithm | How it picks | Good when |
|---|---|---|
| Round robin | Next server in order | Servers equal, requests similar |
| Weighted round robin | Bigger servers get more turns | Mixed server sizes |
| Least connections | Server with fewest open connections | Requests vary in duration |
| Least response time | Fastest recent responder | Latency matters, servers vary |
| IP or key hash | Same client or key goes to same server | Need affinity, such as a per-user cache |
| Random of two choices | Pick two at random, send to the less busy | Large pools; simple and balanced |
Health checks. The balancer regularly calls something like GET /health on each backend. After a few failures it stops sending traffic there; after a few successes it adds the server back. A good health check tests that the server can actually serve (for example, can reach its database), but not so deeply that one shared dependency failing marks every server dead at once.
When to use it. Any time you run more than one copy of a service.
When not to use it. A single internal tool with a few users does not need one. Also, a load balancer does not fix a shared bottleneck behind it, such as one overloaded database.
Common follow-up. "Isn't the load balancer a single point of failure?" Yes, unless you run it redundantly: an active-passive pair sharing a floating IP, several instances behind DNS, or a managed service that is already replicated. Deep dive: load balancing.
Reverse proxy
What it is. A server that sits in front of your application servers and receives client requests on their behalf. Clients talk only to the proxy. (A forward proxy is the opposite: it sits in front of clients and makes requests to the internet on their behalf, as in a corporate network.)
Problem it solves. It gives you one place to handle cross-cutting work: TLS termination (decrypting HTTPS once), compression, caching of responses, serving static files, request size limits, hiding internal server addresses and buffering slow clients so app servers are not tied up.
When to use it. Almost always in front of web applications, even with a single backend, because of the security and convenience benefits.
When not to use it. Do not keep adding proxy layers "just in case"; each hop adds latency and configuration. Many products act as reverse proxy and load balancer at once, so in an interview one box with both responsibilities is fine.
Common follow-up. "What is the difference between a reverse proxy and a load balancer?" A load balancer's job is spreading traffic across several servers; a reverse proxy's job is fronting servers and handling HTTP concerns. A reverse proxy with one backend is still useful; a load balancer with one backend is pointless.
API gateway
What it is. A specialised reverse proxy for APIs that applies policy before requests reach services: authentication, rate limiting, request validation, routing to the right microservice, API version mapping, and sometimes combining several backend calls into one response.
Problem it solves. With many services, you do not want each one to reimplement login checks and quotas, and clients should not know your internal service layout.
When to use it. When you expose APIs to external clients or mobile apps, or when you have several backend services behind one public API.
When not to use it. In a small single-service app, it is extra machinery. Also avoid putting business logic in the gateway; it should stay a thin policy layer, or it becomes a bottleneck that every team must change.
Common follow-up. "Does the gateway replace authorization in services?" No. The gateway can verify who the caller is, but each service must still check whether that caller may access this recipe or order. See identity and security and API design.
Running the code: web servers, app servers and statelessness
Web servers and application servers
What they are. A web server handles HTTP itself: accepting connections, serving static files, passing dynamic requests on. An application server runs your business code: the Recipebox logic that checks permissions, reads data and builds responses. In modern stacks the two are often one process (a Go or Node.js service speaks HTTP directly), with a reverse proxy in front.
Problem they solve. Separating them lets you scale the expensive part (application logic) independently and keep static file serving cheap.
When to separate them. When static traffic is large, or when you want to scale API logic independently. Often you separate further into services: a recipe API, a search API, a notifications worker.
When not to. For small apps, one process behind a managed load balancer is simpler.
Common follow-up. "How many app servers do you need?" Divide peak requests per second by what one server handles at a safe utilisation. The estimation practice lesson works this out.
Stateless services
What it is. A service that stores no per-user data between requests in its own memory or disk. Sessions, uploads and data live in shared stores: a database, a cache, object storage, or a signed token the client carries.
Problem it solves. Stateless servers are interchangeable. Any server can handle any request, so the load balancer can route freely, autoscaling can add or remove machines at any time, and a crash loses nothing.
When to use it. For every web and API tier you want to scale horizontally.
When not to. Some components are inherently stateful: databases, caches, and servers holding long-lived connections such as WebSocket chat servers. You still scale those, but with partitioning and careful routing rather than "add a box".
Common follow-up. "Where do sessions go?" Either in a shared in-memory store keyed by session ID, or in a signed token (such as a signed cookie or JWT) that any server can verify. Tokens avoid a lookup but are hard to revoke before they expire, so keep them short-lived.
Storing data: relational databases and how they scale
Relational database
What it is. A database that stores data in tables of rows and columns, with relationships between tables, queried with SQL. Examples are PostgreSQL and MySQL. Relational databases support ACID transactions: a group of changes either all happen or none do (atomicity), rules stay valid (consistency), concurrent transactions do not see each other's half-done work (isolation) and committed data survives crashes (durability).
Problem it solves. Correctly storing structured data with relationships and flexible queries: "all vegetarian recipes by authors I follow, newest first".
When to use it. By default, for most business data: users, orders, payments, recipes. It is the safest starting point when you are unsure.
When not to. Extremely high write volumes of simple records with no joins (sensor readings, click events), huge unstructured blobs (use object storage), or full-text and fuzzy search (use a search index).
Common follow-up. "Why SQL over NoSQL here?" Because the data is relational, we need transactions (for example, a like and the like count change together) and the scale fits on one primary with replicas for now. Background: SQL fundamentals, transactions and ACID.
Replication: primary-replica and multi-primary
What it is. Keeping copies of the same data on several machines.
- Primary-replica (older texts: master-slave): one primary accepts writes and streams changes to replicas, which serve reads. If the primary fails, a replica is promoted.
- Multi-primary (older texts: master-master): two or more nodes accept writes and replicate to each other.
Problem it solves. Read scaling (more machines to answer reads), availability (a copy survives a machine failure) and locality (a replica near distant users).
| Primary-replica | Multi-primary | |
|---|---|---|
| Writes go to | One node | Any primary |
| Conflicts | None (one writer) | Possible: two primaries change the same row |
| Read scaling | Yes | Yes |
| Write scaling | No | Somewhat, if conflicts are rare |
| Complexity | Moderate (failover, lag) | High (conflict resolution, ordering) |
When to use it. Primary-replica for almost every production database. Multi-primary mainly for multi-region writes where each region must accept writes locally, and even then often with each record "owned" by one region.
When not to. Do not expect replicas to fix a write bottleneck. Do not use multi-primary unless you have a clear conflict strategy, such as last-writer-wins (which silently discards one update) or merging.
Common follow-up. "What is the risk of asynchronous replication?" If the primary confirms a write and crashes before replicas receive it, a promoted replica will not have that write, so it is lost. Synchronous replication to at least one replica prevents this at the cost of write latency. Deep dive: databases at scale.
Federation
What it is. Splitting a database by function: users in one database, recipes in another, comments in a third. Also called functional partitioning.
Problem it solves. Spreads both reads and writes across several databases, and each one is smaller, so it caches better and is easier to back up.
When to use it. When different features have little need to be queried together and one database is overloaded.
When not to. When your main queries join across those areas. Federation turns a join into two queries plus code that combines them, and a single transaction across federated databases needs extra machinery.
Common follow-up. "How do you show a comment with the commenter's name if users and comments are in different databases?" Fetch comments, collect the user IDs, fetch those users in one batch call, and merge in the app. Or store a copy of the display name on the comment (denormalisation) and accept it may lag after a rename.
Sharding
What it is. Splitting one table's rows across several databases ("shards") using a shard key. With hash sharding, the row goes to shard hash(key) mod N. With range sharding, keys in a range (A to F, G to M) go to a given shard. With directory sharding, a lookup table says which shard holds each key.
Problem it solves. One table too large or too write-heavy for any single machine. Each shard holds a fraction of data and receives a fraction of writes.
When to use it. After vertical scaling, replicas, caching and federation are not enough, and only for the tables that need it.
When not to. Early. Sharding makes joins, unique constraints across shards, transactions, reporting and schema changes much harder.
Common follow-up. "What is a hot shard and how do you fix it?" A shard that receives far more traffic than others because of a skewed key, such as one viral recipe. Fixes: choose a key with better spread, split the hot key further (for example recipe_id + random suffix for like counters, summed on read), or cache the hot item. More in scalability.
Denormalisation
What it is. Storing redundant copies of data so reads need fewer joins or lookups. Normalisation (normalization) removes duplication; denormalisation adds some back deliberately.
Problem it solves. Expensive reads. Showing like_count stored on the recipe row is far cheaper than counting rows in a likes table, especially when likes are sharded elsewhere.
When to use it. Read-heavy systems where the same combination is read far more often than written, and after federation or sharding make joins expensive.
When not to. When the copied data changes often or must always be exactly right, unless you can update every copy reliably.
Common follow-up. "How do you keep the copy in sync?" Update it in the same transaction if it is in the same database; otherwise publish an event and let a worker update copies, accepting a short delay. Periodically reconcile counts from the source of truth.
NoSQL databases: four main types
NoSQL is an umbrella term for databases that do not use the relational table model. Most give up some query flexibility or transaction scope in exchange for easier horizontal scaling or a data model that fits a specific access pattern. Many offer BASE semantics (basically available, soft state, eventually consistent) rather than full ACID across all data, though several now support transactions within limits. Deep dive: NoSQL and distributed databases.
| Type | Data model | Good for | Weak at | Typical examples |
|---|---|---|---|---|
| Key-value | Value looked up by a key | Sessions, caches, counters, user settings | Querying by anything but the key | Redis, DynamoDB (also document) |
| Document | JSON-like documents, each with its own fields | Product catalogues, user profiles, content with varying fields | Joins across documents | MongoDB, Couchbase |
| Wide-column | Rows with many columns, grouped by partition key, sorted within | Time series, event logs, very high write volume | Ad hoc queries; must design tables per query | Cassandra, HBase, Bigtable |
| Graph | Nodes and edges with properties | Relationship-heavy queries: friends of friends, recommendations, fraud rings | Bulk analytics over everything; sharding graphs is hard | Neo4j, Amazon Neptune |
When to choose NoSQL over SQL. When access patterns are simple and known (look up by key, list by partition), data volume or write rate is beyond one relational primary, the schema varies a lot between records, or the data is naturally a graph.
When to stay with SQL. When you need transactions across many records, flexible queries you cannot predict, or strong constraints such as uniqueness and foreign keys.
Common follow-up. "Your chat app stores messages in a wide-column store. What is the partition key?" Usually conversation_id, with messages sorted by time inside the partition, so "latest 50 messages in this conversation" reads one partition in order. If one conversation could grow forever, add a time bucket to the key, such as (conversation_id, month), to keep partitions bounded.
Caching: layers and patterns
Cache layers
What it is. A cache is a fast, smaller store holding copies of data that is expensive to fetch or compute. Caches exist at many layers, each closer to the user than the last:
Browser cache --> CDN edge --> Reverse proxy cache
--> App in-process cache --> Distributed cache --> DB buffer cache
- Client (browser) cache: controlled by HTTP headers such as
Cache-Control. - CDN: shared cache for public content.
- Web server or reverse proxy cache: caches whole HTTP responses.
- Application in-process cache: a small map inside each app server; fastest, but each server has its own copy.
- Distributed cache: a shared in-memory cluster such as Redis or Memcached, reached over the network.
- Database cache: the database keeps hot pages in memory itself.
Problem it solves. Lower latency and less load on slower systems.
When to use it. Read-heavy data that is read many times between changes and where slight staleness is acceptable.
When not to. Data that changes on almost every read, data that must be exactly current (an account balance before a withdrawal), or when the hit ratio would be low because requests rarely repeat.
What to cache. Either database query results (key = query) or application objects (key = recipe:42, value = the fully assembled recipe). Objects are usually easier: when recipe 42 changes, you delete one key, whereas with query caching you may not know which cached queries included it.
The four cache patterns
| Pattern | Read path | Write path | Strength | Weakness |
|---|---|---|---|---|
| Cache-aside (lazy loading) | App checks cache; on miss reads DB and fills cache | App writes DB, then deletes cache key | Only requested data is cached; cache failure is not fatal | First read is slow; stale window if delete fails |
| Write-through | App reads cache | App writes cache, which synchronously writes DB | Cache always fresh for written data | Slower writes; caches data nobody reads |
| Write-behind (write-back) | App reads cache | App writes cache; cache writes DB later, in batches | Very fast writes; batched DB load | Data loss if cache dies before flushing; complex |
| Refresh-ahead | App reads cache | Cache reloads popular keys shortly before they expire | Hot keys rarely miss | Wasted work if prediction is wrong |
Here is cache-aside in Python-style pseudocode:
def get_recipe(recipe_id):
key = f"recipe:{recipe_id}"
cached = cache.get(key)
if cached is not None:
return cached # hit
recipe = db.query_one("SELECT * FROM recipes WHERE id = %s", recipe_id)
cache.set(key, recipe, ttl_seconds=600) # fill on miss
return recipe
def update_recipe(recipe_id, fields):
db.update("recipes", recipe_id, fields)
cache.delete(f"recipe:{recipe_id}") # next read refills
Note that the write path deletes the key rather than writing the new value. If two updates race, setting values can leave the older one in the cache; deleting is safer.
Eviction. When the cache is full it must drop something. LRU (least recently used) is the common default; LFU (least frequently used) keeps long-term favourites; TTL expiry removes entries after a fixed time regardless.
Common follow-up. "What is a cache stampede and how do you prevent it?" When a hot key expires, thousands of requests miss at once and all hit the database. Prevent it by letting only one request rebuild the key (a lock or "single flight"), serving the stale value briefly while it refreshes, or adding random jitter to TTLs so keys do not expire together. Deep dive: caching.
Asynchronous work: queues and back pressure
Message queues vs task queues
What they are. A message queue holds messages from producers until consumers read them, so the two sides do not need to be running at the same moment or at the same speed. A task queue is a message queue specialised for jobs: each message names a function to run and its arguments, and a pool of workers executes them, often with retries, scheduling and result tracking.
| Message queue | Task queue | |
|---|---|---|
| A message means | "Something happened" or "here is data" | "Please do this job" |
| Typical consumer | Any service interested in the event | A worker pool running a known function |
| Examples | RabbitMQ, Amazon SQS, Kafka (a log, similar role) | Celery, Sidekiq, background job frameworks |
| Recipebox use | "Recipe 42 created" consumed by search and notifications | "Resize photo 42 to 3 sizes" |
Problem they solve. Slow, unreliable or bursty work no longer blocks the user's request; producers and consumers scale independently; failed work is retried.
When to use them. Emails, notifications, image and video processing, search indexing, webhooks to third parties, anything that can finish a few seconds later.
When not to. When the user needs the result in the response (checking a password) or strict immediate consistency is required.
Common follow-up. "What delivery guarantee do you get?" Most queues give at-least-once delivery: a message is redelivered if the consumer crashes before acknowledging it, so consumers must be idempotent (processing twice has the same effect as once), often by recording processed message IDs. Exactly-once processing is achievable only within limits, usually by combining at-least-once delivery with idempotent writes. See message queues.
Back pressure
What it is. A way for an overloaded component to tell upstream components to slow down, instead of accepting work it cannot finish.
Problem it solves. Without it, a queue grows without limit, messages get hours old, memory runs out and the system falls over. Accepting work you cannot do is worse than refusing it quickly.
How it looks in practice. Bounded queues that reject new items when full; returning HTTP 429 Too Many Requests or 503 Service Unavailable with a Retry-After header; clients retrying with exponential backoff (wait 1 s, 2 s, 4 s, plus random jitter).
When to use it. Anywhere a fast producer feeds a slower consumer.
Common follow-up. "Why add jitter to retries?" If thousands of clients fail at the same moment and all retry after exactly 2 seconds, they hit the server together again. Random jitter spreads the retries out.
Files and search
Object storage
What it is. A service that stores whole files ("objects") under keys and serves them over HTTP, such as Amazon S3, Google Cloud Storage or Azure Blob Storage. It is durable, by keeping several copies, and scales to effectively unlimited size.
Problem it solves. Storing images, videos, backups and documents cheaply and durably without filling app server disks or bloating databases.
When to use it. Any large binary data: photos, videos, uploads, exports, logs, backups. Put a CDN in front for public reads; store only the key in your database.
When not to. Small frequently-updated records, data you need to query by fields, or files you need to edit in place (objects are replaced whole).
Common follow-up. "How do clients upload large files without going through your servers?" The server checks permissions and returns a pre-signed URL, a short-lived link that allows one upload to one key; the client uploads directly to object storage, and the app is notified when it completes. See storage, search and geo and the file sync case study.
Search index
What it is. A separate data store optimised for text search, typically built on an inverted index: a map from each word to the list of documents containing it, like the index at the back of a book. Examples include Elasticsearch and OpenSearch.
Problem it solves. Queries like "recipes containing paneer and spinach, tolerant of typos, ranked by relevance" are slow or impossible with SQL LIKE '%paneer%', which scans every row.
When to use it. Full-text search, autocomplete, faceted filtering ("vegetarian, under 30 minutes").
When not to. As your source of truth. The index is a derived copy, updated asynchronously from the database, and it can be rebuilt.
Common follow-up. "How do you keep the index in sync?" Publish change events (or use change data capture, which reads the database's change log) and have a worker update the index. Expect a short delay before new recipes appear in search. See transactions and CDC.
Distribution helpers
Consistent hashing
What it is. A way to assign keys to servers so that adding or removing a server moves only a small fraction of keys. Imagine a circle of hash values. Each server is placed at one or more points on the circle; each key is hashed onto the circle and belongs to the next server clockwise.
0
S3 . . S1
. k1 . k1 -> next server clockwise = S1
. . k2 -> S2
. k2 .
S2 . . add S4 between S1 and S2:
| only keys between S1 and S4 move
Problem it solves. With hash(key) mod N, changing N from 4 to 5 moves roughly 80% of keys, causing a flood of cache misses or data movement. With consistent hashing, adding the fifth server moves only about one fifth of the keys, the share the new server takes over.
Virtual nodes. Each physical server is placed at many points (for example 100) on the circle, which evens out the load and lets a bigger server take more points.
When to use it. Distributed caches, partitioned key-value stores, sharded databases where nodes join and leave.
When not to. When the number of shards is fixed and rarely changes; a fixed mapping or a directory table may be simpler.
Common follow-up. "Does consistent hashing fix hot keys?" No. It balances many keys across servers, but one extremely popular key still lands on one server. You need replication of that key or key splitting. More in scalability.
Rate limiter
What it is. A component that caps how many requests a client (by user ID, API key or IP) may make in a time window, returning 429 Too Many Requests beyond the limit.
Problem it solves. Protects services from abuse, runaway scripts and accidental overload, and enforces fair use between customers.
Common algorithms: token bucket (a bucket refills at a steady rate; each request spends a token; allows short bursts), leaky bucket (requests drain at a fixed rate; smooths bursts), fixed window counter (count per calendar minute; simple but allows double bursts at window edges) and sliding window (counts over the last 60 seconds; more accurate).
When to use it. Public APIs, login and OTP endpoints, expensive operations such as AI calls or search.
When not to. Do not use it as your only capacity plan. Rate limits protect against unfair use; they do not replace provisioning for legitimate peak load.
Common follow-up. "How do you rate-limit across 20 API servers?" Keep counters in a shared fast store (such as Redis, with atomic increments and expiry), or give each server a share of the limit and accept slight inaccuracy. Implementation detail: LRU cache and rate limiter design.
Service discovery
What it is. A way for services to find the current network addresses of other services. Instances register themselves (or are registered by the platform) in a registry, and callers look them up by name. Examples: Consul, etcd, ZooKeeper, and the built-in service DNS in Kubernetes.
Problem it solves. With autoscaling and failures, instance addresses change constantly. Hard-coding IPs breaks.
When to use it. Microservices on dynamic infrastructure.
When not to. A monolith behind one load balancer with a stable DNS name needs nothing more.
Common follow-up. "Client-side vs server-side discovery?" In client-side discovery the caller fetches the list of instances and load-balances itself; in server-side discovery the caller hits a load balancer or DNS name that resolves to healthy instances. The registry must itself be highly available, which is why it often runs on a consensus system; see consensus and coordination.
Monitoring
What it is. Collecting and acting on signals about system health: metrics (numbers over time such as request rate, error rate, latency percentiles, CPU), logs (records of individual events) and traces (the path of one request across services). Alerts notify people when something is wrong.
Problem it solves. You cannot fix or scale what you cannot see. Monitoring finds the next bottleneck and tells you about outages before users do.
A simple starting set. For every service, watch the "four golden signals": latency, traffic, errors and saturation (how full the most constrained resource is). Add component-specific ones: cache hit ratio, queue depth, replication lag, disk space.
When to use it. From day one.
Common follow-up. "Why alert on p99 latency rather than average?" Averages hide slow outliers: if 1% of requests take 5 seconds, the average barely moves, but 1 in 100 users has a terrible experience, and a page that makes 50 calls almost always hits one slow call. See observability and performance.
Consistency patterns
When data is copied (replicas, caches, regions), a read may hit a copy that has not yet received the latest write. A consistency pattern describes what readers are guaranteed to see.
| Pattern | Guarantee | Example in Recipebox | Cost |
|---|---|---|---|
| Weak | A read may or may not see a write; no promise | Live "people cooking now" counter; a dropped update does not matter | Cheapest, fastest |
| Eventual | If writes stop, all copies converge to the same value, usually within milliseconds to seconds | Like counts, feeds, search results | Brief staleness; must design UI for it |
| Strong | Every read after a write completes sees that write | Username claims, payments, account deletion | Higher latency; lower availability during partitions |
Weak consistency is common in real-time systems where a missed update is soon replaced, such as live video calls or game positions. Eventual consistency is the default for replicated and cached data at scale. Strong consistency usually means synchronous replication or reading through a single leader, sometimes with a quorum (a majority of replicas agreeing).
There are useful middle grounds, covered in fundamentals: read-your-own-writes (you always see your own changes), monotonic reads (you never see data go backwards in time) and causal consistency (if one event caused another, everyone sees them in that order).
Common mistake
Saying "eventual consistency means data might be wrong forever". It means copies converge once writes stop; the question is how long the window is and what the user sees during it. Always state the expected delay and how the UI handles it.
Availability patterns
Fail-over: active-passive and active-active
Active-passive (also called primary-standby): one node serves traffic; a standby watches it through heartbeats (periodic "I am alive" messages). If heartbeats stop, the standby takes over the address and begins serving. A hot standby is already running and synced, so failover takes seconds; a cold standby must be started and loaded, so failover takes minutes or more.
Active-active: all nodes serve traffic at once, and a load balancer or DNS spreads requests. If one fails, the others absorb its share, which means each must have spare capacity.
| Active-passive | Active-active | |
|---|---|---|
| Normal capacity used | Half (standby idles) | All |
| Failover time | Seconds to minutes | Near zero for stateless tiers |
| Data complexity | One writer, simpler | Concurrent writers for databases, harder |
| Risk | Standby untested, fails when needed | Overload if survivors lack headroom |
Fail-over risks: data written to the old primary but not replicated may be lost, and both nodes may believe they are primary (split brain), which is why systems use consensus or "fencing" to ensure only one leader writes. See reliability and recovery.
Replication as an availability pattern
Replication (described above) keeps a copy of data ready on another machine, so losing one machine does not lose data or service. Fail-over without replicated data is useless; replication is what makes fail-over possible for stateful tiers.
Availability in numbers
Availability is usually stated as the percentage of time a service works, written as "nines". The allowed downtime for each level:
| Availability | Downtime per year | Per month (30 days) | Per day |
|---|---|---|---|
| 99% ("two nines") | 87.6 hours | 7.2 hours | 14.4 minutes |
| 99.9% ("three nines") | 8.76 hours | 43.2 minutes | 86.4 seconds |
| 99.99% ("four nines") | 52.6 minutes | 4.32 minutes | 8.64 seconds |
| 99.999% ("five nines") | 5.26 minutes | 25.9 seconds | 0.86 seconds |
How to compute it: downtime = (1 - availability) times the period. For 99.9% over a year, 0.001 times 525,600 minutes = 525.6 minutes, which is 8.76 hours.
Components in sequence. If a request must pass through several components, all must work, so you multiply their availabilities:
Total = A1 x A2 x ... x An
Components in parallel. If either of two redundant copies can serve, the system fails only when both fail:
Total = 1 - (1 - A1) x (1 - A2)
Worked example: Recipebox's request path
A Recipebox API request passes through a load balancer (99.99%), one app server (99.9%) and one database (99.9%). Assume failures are independent.
Step 1, all in sequence:
0.9999 x 0.999 x 0.999 = 0.99790 (about 99.79%)
Downtime per year = 0.00210 x 525,600 min = about 1,103 min
= about 18.4 hours
Notice the result is worse than the weakest component. Chaining reduces availability.
Step 2, two app servers in parallel:
App tier = 1 - (0.001 x 0.001) = 0.999999
Total = 0.9999 x 0.999999 x 0.999 = 0.99890 (about 99.89%)
Downtime per year = about 579 min = about 9.6 hours
Better, but now the single database dominates.
Step 3, add a database replica with automatic failover (treated as parallel):
DB tier = 1 - (0.001 x 0.001) = 0.999999
Total = 0.9999 x 0.999999 x 0.999999 = 0.99990 (about 99.99%)
Downtime per year = about 54 min
Now the load balancer's 99.99% is the limit. The lesson: redundancy at every tier is what moves you up the nines, and the weakest non-redundant component sets the ceiling.
Common mistake
The parallel formula assumes failures are independent. Two app servers in the same rack, with the same bad deploy, or behind the same failed switch fail together. Real availability is lower than the formula unless copies are truly separated (different machines, zones and rollout waves).
Communication: TCP vs UDP, RPC vs REST
TCP vs UDP
TCP (Transmission Control Protocol) sets up a connection with a handshake, then guarantees bytes arrive complete and in order, retransmitting lost packets and slowing down when the network is congested. UDP (User Datagram Protocol) sends individual packets ("datagrams") with no connection, no ordering and no retransmission; the application decides what to do about loss.
| TCP | UDP | |
|---|---|---|
| Reliability and ordering | Built in | Not built in (app can add it) |
| Setup | Handshake before data | None |
| Head-of-line blocking | Yes: a lost packet delays everything after it | No |
| Use for | Web pages, APIs, databases, file transfer, email | Live voice and video, online games, DNS queries, QUIC (HTTP/3) |
Rule of thumb: use TCP when every byte must arrive correctly; use UDP when a late packet is useless anyway (a video frame from 300 ms ago) and low latency matters more. Details: transport layer: TCP and UDP.
RPC vs REST
RPC (remote procedure call) makes calling another service look like calling a function: GetRecipe(id=42). Frameworks such as gRPC define the functions and message types in a schema file and generate client code; gRPC uses HTTP/2 and a compact binary format (Protocol Buffers).
REST (representational state transfer) models the API as resources with URLs, manipulated with standard HTTP methods: GET /recipes/42, POST /recipes, DELETE /recipes/42. Responses are usually JSON and can use standard HTTP caching.
| RPC (for example gRPC) | REST | |
|---|---|---|
| Model | Actions (functions) | Resources (nouns) with HTTP verbs |
| Payload | Often binary, schema-defined | Usually JSON, human-readable |
| Strengths | Fast, strongly typed, streaming, good for internal service-to-service calls | Simple, universal, cacheable, easy to debug with a browser or curl |
| Weaknesses | Harder for browsers and third parties; tighter coupling to the schema | More verbose; some actions do not map naturally to resources |
| Typical use | Internal microservice calls | Public APIs, web and mobile clients |
A common real-world split: REST (or GraphQL) at the public edge, gRPC between internal services. See API design.
One-table cheat sheet
| Block | One-line job | Reach for it when | Watch out for |
|---|---|---|---|
| DNS | Name to address, coarse routing | Always; region steering | Cached answers slow failover |
| CDN | Serve content near users | Static or public content | Stale files, caching private data |
| Load balancer | Spread traffic, skip dead servers | More than one server | Itself a SPOF; hides DB bottleneck |
| Reverse proxy | Front door for HTTP concerns | Any web app | Extra hop if overused |
| API gateway | Auth, quotas, routing for APIs | Many services or public API | Business logic creeping in |
| Stateless app servers | Interchangeable compute | Horizontal scaling | Sessions and files must move out |
| Relational DB | Structured data with transactions | Default for business data | Write ceiling of one primary |
| Replication | Copies for reads and availability | Read-heavy, need failover | Lag; writes do not scale |
| Federation | Split DB by feature | Separate features overloaded | No cross-DB joins |
| Sharding | Split table by key | One table too big or write-heavy | Hot shards, resharding, no joins |
| Denormalisation | Store copies to avoid joins | Read-heavy, joins expensive | Keeping copies in sync |
| Key-value store | Fast lookup by key | Sessions, counters, caches | No queries beyond the key |
| Document store | Flexible JSON records | Varying fields per record | Weak joins |
| Wide-column store | Huge write volumes, partitioned | Time series, logs, messages | Must design per query |
| Graph DB | Relationship traversal | Friends-of-friends, fraud | Hard to shard |
| Cache | Fast copy of hot data | Read-heavy, staleness OK | Invalidation, stampedes |
| Message / task queue | Do work later, decouple | Slow or flaky side effects | Duplicates; need idempotency |
| Back pressure | Refuse what you cannot finish | Fast producer, slow consumer | Clients must retry with backoff |
| Object storage | Durable cheap file storage | Images, video, backups | Not queryable, whole-object writes |
| Search index | Text and faceted search | Search, autocomplete | Derived copy, lags the DB |
| Consistent hashing | Move few keys when nodes change | Caches, partitioned stores | Does not fix one hot key |
| Rate limiter | Cap requests per client | Public APIs, logins | Not a capacity plan |
| Service discovery | Find live instances by name | Dynamic microservices | Registry must be highly available |
| Monitoring | See health, find bottlenecks | Always | Alert on symptoms users feel |
Interview questions
Q1. What is the difference between latency and throughput, and can improving one hurt the other?
Latency is the time for one request; throughput is requests completed per unit time. Yes: batching raises throughput but adds waiting time per item, and running servers near 100% utilisation maximises throughput while queues grow and latency rises sharply. You usually maximise throughput subject to a latency target such as p99 under 300 ms.
Q2. L4 vs L7 load balancer?
An L4 balancer routes by IP and port per connection without reading the content, which is fast and works for any TCP or UDP protocol. An L7 balancer understands HTTP, so it can route by path or header, terminate TLS, retry and add headers, at higher cost per request. Use L7 for web APIs needing smart routing, L4 for raw throughput or non-HTTP protocols.
Q3. Explain cache-aside and why writes delete the key instead of updating it.
The app reads from the cache, and on a miss loads from the database and fills the cache. On writes it updates the database and deletes the key, so the next read reloads fresh data. Deleting avoids a race where two concurrent updates write their values to the cache in the wrong order and leave stale data indefinitely.
Q4. When would you choose write-behind caching, and what is the risk?
When writes are very frequent and can be batched, such as view counters, because the app writes only to memory and the cache flushes to the database periodically. The risk is losing recent writes if the cache node fails before flushing, so it suits data where small losses are acceptable or the cache itself is replicated and persisted.
Q5. Message queue vs task queue?
A message queue carries messages or events between producers and consumers, decoupling them in time and speed; consumers decide what to do. A task queue is a specialisation where each message is a job to execute, with worker pools, retries and scheduling built in. Both usually deliver at least once, so consumers must be idempotent.
Q6. What is back pressure and how would you implement it for an upload pipeline?
Back pressure is signalling upstream to slow down when a component is saturated. For uploads, cap the processing queue length; when it is full, the API returns 503 or 429 with Retry-After, and clients retry with exponential backoff and jitter. Monitoring queue depth and age lets you add workers before limits are hit.
Q7. Why does consistent hashing matter for a distributed cache?
With modulo hashing, changing the number of cache nodes remaps most keys, causing a burst of misses that can overload the database. Consistent hashing remaps only the keys owned by the added or removed node, roughly 1/N of them. Virtual nodes keep the distribution even.
Q8. Strong vs eventual consistency: give one feature for each.
Strong consistency suits operations where a stale read causes real harm: claiming a unique username, deducting a wallet balance, or revoking access. Eventual consistency suits counts, feeds and search results, where a few seconds of delay is acceptable and the higher availability and lower latency are worth it.
Q9. Compute the availability of two 99.9% components in series and in parallel.
In series both must work, so 0.999 times 0.999 = 0.998 (99.8%), about 17.5 hours of downtime a year. In parallel the system fails only if both fail: 1 minus 0.001 times 0.001 = 0.999999 (99.9999%). This assumes independent failures, which real systems only approximate.
Q10. Active-passive vs active-active failover?
Active-passive keeps a standby that takes over when heartbeats from the active node stop; it is simpler but idles capacity and failover takes time. Active-active serves from all nodes, so failure just shifts load, but each node needs headroom and stateful data with multiple writers needs conflict handling. Stateless tiers are usually active-active; databases are often active-passive.
Q11. SQL or NoSQL for a user profile service with 500 million users?
Profiles are looked up by user ID, have some optional fields and are read far more than written, which fits a key-value or document store well. A relational database with sharding by user ID also works and keeps constraints like unique email easier. I would choose based on whether we need cross-user queries and transactions; if not, a horizontally scalable document store is simpler to operate at that size.
Q12. RPC or REST between internal microservices?
gRPC is common internally: a typed schema, generated clients, compact binary messages and HTTP/2 streaming reduce bugs and latency. REST is simpler to debug and universally supported, so it remains common at the public edge. Either works; consistency within the organisation matters more than the choice.
Q13. Why use object storage plus a CDN for images instead of storing them in the database?
Databases are optimised for small structured rows; storing large blobs bloats backups, replication and memory. Object storage is cheap, durable and serves files over HTTP, and a CDN in front caches them near users. The database keeps only the object key and metadata.
Q14. What would you monitor first for a new service?
The four golden signals: request rate, error rate, latency percentiles (p50, p95, p99) and saturation of the tightest resource. Then component-specific metrics such as cache hit ratio, queue depth and replication lag. Alerts should fire on what users feel, such as error rate or p99 latency, not on every CPU spike.
Key takeaways
- Every building block answers a specific problem; know its "when not to use" as well as its benefit.
- Performance is speed for one user; scalability is keeping that speed as load grows.
- Replicas and caches scale reads; federation and sharding scale writes.
- Cache-aside with delete-on-write is the safe default; know write-through, write-behind and refresh-ahead trade-offs.
- Queues decouple slow work but bring at-least-once delivery, so consumers must be idempotent; back pressure stops them overflowing.
- Choose consistency per feature: strong for money and uniqueness, eventual for counts and feeds.
- Availability multiplies in series and improves with independent redundancy; the weakest single component sets the ceiling.
- TCP for correctness, UDP for real-time; REST at the public edge, RPC often inside.
Next lesson
Continue with Estimation practice.

