50 System Design Concepts, Explained Properly
Fifty system design concepts explained in plain English - what each one is, why it exists, what problem it solves, and what it costs. Covering infrastructure, data and storage, distributed systems, messaging, reliability, and the AI concepts that have become standard vocabulary.
Published on 20 sept 2026

There is a particular kind of quiet panic that shows up in design reviews.
Someone says a word you have definitely heard before - backpressure, idempotency, consistent hashing - and you nod, because you recognise it, and because asking would mean admitting that recognising a word and understanding it are not the same thing. The conversation moves on. You look it up later, or you don't.
System design has more of these moments than almost any other part of engineering. The vocabulary is large, the ideas are abstract, and almost nobody explains them from the beginning, because everyone in the room is assumed to know them already. So you end up with engineers - good ones - who can use these words correctly in a sentence but could not tell you why the thing exists or what it costs.
This is a guide to fifty of those words. Not dictionary definitions. For each one: what it is, why it exists, what problem it solves, and where the trade-off lives. That last part matters most. A concept you can define is trivia. A concept whose cost you understand is something you can actually make a decision with.
How to read this. You do not need to go front to back. The five sections are roughly independent, and each concept stands on its own. If you are preparing for an interview, the trade-off line at the end of each one is the part worth remembering - it is almost always the follow-up question.
Section 1: Core Infrastructure
How traffic reaches your system, how it gets spread out, and how the whole thing grows.
1. Scalability
Scalability is a system's ability to handle growth without falling over. Growth might mean more users, more data, more requests per second, or users spread across more of the world.
The part people miss: scalability is not a property of any one component. It belongs to the whole system, and the system is only as scalable as its worst part. You can scale your application servers to handle ten times the traffic and discover that all you did was move the queue - now the database is the thing falling over. Every layer has to be able to grow, or the bottleneck just relocates.
The trade-off: everything below. Scalability is bought with complexity, money, or both.
2. Vertical Scaling
Vertical scaling means making one machine bigger. More CPU, more memory, faster disks. The application does not change. The architecture does not change. You buy a larger box.
It is genuinely the simplest option and it works well right up until it doesn't. Two limits: machines have a maximum size, and one machine is one point of failure. There is no redundancy in a single server. When it goes down, all of it goes down.
The trade-off: simplicity now, a hard ceiling and no redundancy later.
3. Horizontal Scaling
Horizontal scaling means adding more machines rather than bigger ones. Need more capacity, run more copies. There is no ceiling, and one machine dying costs you a slice of capacity instead of the whole service.
The catch is that the servers have to be stateless - holding nothing between requests that belongs to a particular user. The moment a server remembers something about you, your next request has to come back to that same server, and you have lost the flexibility that made horizontal scaling work in the first place.
The trade-off: no ceiling, but your application has to be built for it, and state has to live somewhere else.
4. Load Balancer
A load balancer sits in front of a group of servers and shares the incoming requests between them. Every request hits the balancer first; it picks a server, passes the request on, and returns the answer. From outside it looks like one server. Inside, many are splitting the work.
┌──────────────┐ users ─────────▶ │ load balancer│ └──────┬───────┘ ┌─────────┼─────────┐ ▼ ▼ ▼ ┌─────┐ ┌─────┐ ┌─────┐ │ srv │ │ srv │ │ srv │ ✗ unhealthy │ 1 │ │ 2 │ │ 3 │ ── removed from └─────┘ └─────┘ └─────┘ rotation
It also keeps checking whether each server behind it is healthy, and stops sending traffic to ones that are failing. That automatic noticing-and-rerouting is what turns a pile of servers into something that survives one of them dying.
The trade-off: the balancer itself must not become the single point of failure you just eliminated.
5. Latency
Latency is the time between sending a request and getting the answer - the delay one user feels. Low latency feels instant. High latency feels like the page is broken.
Latency comes from several places stacked on top of each other: the time the request spends travelling over the network, the time your server spends thinking, the time the database spends answering. Reducing it means finding which of those dominates and fixing that one. Optimising the wrong layer is the most common wasted week in performance work.
The trade-off: many latency fixes (caching, replicas, CDNs) buy speed with staleness.
6. Throughput
Throughput is how much work the system finishes in a period of time - usually requests per second. It describes total capacity, not the speed of any single request.
Latency and throughput are related but not the same, and confusing them causes bad decisions. A system can be quick per request and still have terrible throughput if it handles one at a time. A system can be slow per request and have excellent throughput if it handles thousands in parallel. Designing for one does not give you the other for free.
The trade-off: batching and parallelism raise throughput, often at the cost of per-request latency.
7. CDN (Content Delivery Network)
A CDN is a fleet of servers spread around the world that keep copies of your static files and serve them from wherever is closest to each user. Instead of everyone fetching the same image from one machine on another continent, they fetch it from a machine in their own city.
Two things improve. Latency drops, because the bytes travel a shorter distance. And your origin server sheds an enormous amount of traffic it no longer has to serve.
The trade-off: cached copies go stale. When the file changes, you need the cache cleared quickly, and "quickly" across hundreds of locations is its own problem.
8. DNS (Domain Name System)
DNS turns a name a person can remember, like example.com, into the IP address a computer needs to actually connect. It is the internet's phone book.
It shows up in system design because it can route traffic: send users to the nearest region, or steer them away from a datacentre that is unhealthy. It is the simplest form of global routing you can have.
The trade-off: DNS answers get cached everywhere - by browsers, by operating systems, by resolvers in between - so changes take minutes to hours to take effect. It is not a mechanism for fast failover.
9. API Gateway
An API gateway is one front door sitting in front of many backend services. It handles the things every service needs but none of them should each implement: checking who the caller is, enforcing rate limits, terminating TLS, and routing the request to the right service.
┌──────────────────────────┐ client ────────▶ │ API gateway │ │ auth · rate limit · TLS │ └────┬─────────┬───────┬───┘ ▼ ▼ ▼ orders payments users service service service
Without one, every service writes its own authentication, and they drift apart, and one of them gets it wrong. With one, the concern lives in a single place.
The trade-off: the gateway is now on the critical path for everything. It has to be highly available and fast, or it becomes the bottleneck for your whole platform.
10. Reverse Proxy
A reverse proxy takes requests on behalf of one or more servers and passes them along. The client talks to the proxy and neither knows nor needs to know which machine actually did the work.
Reverse proxies do caching, compression, TLS, and load balancing. The relationship to a load balancer confuses people, so: a load balancer is a reverse proxy that specialises in spreading traffic. A reverse proxy is the broader category and can do plenty of other jobs.
The trade-off: another hop, and another thing to operate.
Section 2: Data and Storage
Where the data lives, how it is found, and how it stays correct as the system grows.
11. Database
A database stores, organises, and retrieves data that has to outlive a single request. Without one, everything created while handling a request vanishes the moment the request ends.
Databases come in many shapes, each suited to different kinds of data and different ways of reading it. This is among the most consequential choices in any design, because it is painful to change later and everything else ends up depending on it.
The trade-off: you are picking which access patterns will be cheap and which will be expensive, often before you know what they are.
12. SQL Database
A SQL database - also called relational - stores data in tables of rows and columns, and enforces the relationships between those tables. You query it with a structured language, and it gives you ACID guarantees (see next).
It is the right answer when the data is structured, when things relate to other things, when you need transactions, and when you do not yet know every query you will want to run. That last point is underrated: relational databases are forgiving of questions you did not plan for.
The trade-off: scaling writes past a single machine means sharding, and sharding is a significant step up in complexity.
13. NoSQL Database
NoSQL is an umbrella over everything that is not a relational table: key-value stores, document stores, wide-column stores, graph databases. Each one gives up some of SQL's query flexibility or consistency in exchange for scale, flexibility, or speed on one particular access pattern.
The usual mistake is choosing NoSQL because it sounds current rather than because the data actually fits. The way you will read the data should decide this, not the technology's reputation.
The trade-off: you optimise hard for the access patterns you designed around, and pay dearly for the ones you did not.
14. ACID
ACID is four guarantees a database transaction gives you:
| Letter | Means | In practice |
|---|---|---|
| Atomicity | All or nothing | Every step in the transaction succeeds, or none of them do |
| Consistency | Valid state to valid state | The database never lands in a state its rules forbid |
| Isolation | Transactions don't see each other mid-flight | Concurrent work doesn't interfere |
| Durability | Committed means committed | Once it says done, a crash won't lose it |
This is exactly why money lives in relational databases. For a bank transfer, half-finished is worse than not started - atomicity is not a nice-to-have, it is the whole point.
The trade-off: these guarantees cost coordination, and coordination costs latency.
15. Index
An index is a second data structure the database maintains so that certain lookups are fast. Without one, finding the rows matching a condition means reading every row in the table. With one, the database jumps more or less straight to them.
The cost is on the other side. Indexes make reads faster and writes slower, because every insert, update, and delete has to update every index on that table. Ten indexes means paying ten times the index-maintenance cost on every write.
The trade-off: index the queries you actually run. Not every column, and not the ones you imagine you might run one day.
16. Sharding
Sharding splits your data across several databases so each one holds only a slice. A shard key decides which slice a given row belongs to. Because the load is divided, you scale both write throughput and total storage.
┌──────────── shard key: user_id ────────────┐ ▼ ▼ ▼ ┌──────────┐ ┌──────────┐ ┌──────────┐ │ shard A │ │ shard B │ │ shard C │ │ users │ │ users │ │ users │ │ 0 – 33 M │ │ 33 – 66 M│ │ 66 – 99 M│ └──────────┘ └──────────┘ └──────────┘ a query for one user → one shard (fast) a query across users → all shards (slow, must be merged)
The hard part is choosing the key. A good one spreads both the data and the traffic evenly. A bad one creates a hot partition - one shard taking most of the load while the rest idle. And any query that does not include the shard key has to run on every shard and have its results stitched together.
The trade-off: you gain write scale and lose easy cross-cutting queries, plus you now have a key you cannot change without moving all your data.
17. Replication
Replication keeps copies of your data on more than one machine. The usual arrangement is leader-follower: one primary takes all the writes, and replicas hold copies that serve reads. You get read capacity and redundancy from the same mechanism.
The subtlety worth knowing is replication lag - the gap between a write landing on the primary and appearing on the replicas. During that window, a read from a replica returns the old value. Usually that is fine. Occasionally it is a bug report titled "I saved it and it didn't save", and you need to route that particular read to the primary.
The trade-off: more read capacity, in exchange for reads that are sometimes slightly behind reality.
18. Cache
A cache is a fast layer holding copies of data you ask for often, so you can skip the slow work of fetching it again. The classic version is an in-memory cache in front of a database, answering reads from memory and only troubling the database when it does not have the answer.
The cost is always the same one: the cache holds a copy, and copies go out of date when the original changes. Designing a cache is really deciding how stale is acceptable and for how long - which is what your expiry and invalidation rules encode.
The trade-off: speed for staleness. There is no version of caching that avoids this.
19. Cache-Aside
Cache-aside is the most common caching pattern, and the sensible default. The application looks in the cache first. On a hit, it returns what it found. On a miss, it reads the database, puts the result in the cache, and returns it. The cache fills up lazily, with exactly the data people actually ask for.
read ──▶ cache? ──hit──▶ return value │ miss │ ▼ read database ──▶ store in cache ──▶ return value
Two things make it a good default: the cache naturally ends up holding the hot data rather than whatever you guessed would be hot, and if the cache goes down entirely, the application just gets slower instead of breaking.
The trade-off: the first request for anything is always slow, and two simultaneous misses will both hit the database.
20. Write-Through
Write-through means every write goes to the cache and the database together. The cache is never out of step with the database, because nothing gets written to one without the other.
You pay for that in write latency: every write waits for both to finish before it returns. It is the right call when a user has to see their own write immediately afterwards and the extra milliseconds are acceptable.
The trade-off: consistency bought with slower writes.
21. Write-Behind
Write-behind sends writes to the cache, returns immediately, and flushes them to the database later in the background. Writes feel instant because the cache acknowledges them straight away.
The risk is plain: if the cache dies before it flushes, those writes are gone. That makes this appropriate for high-volume work where losing a little is survivable - analytics events, view counters - and inappropriate for anything you would have to explain to a customer.
The trade-off: the fastest writes of the three patterns, paid for with a real window of possible data loss.
22. Consistent Hashing
Consistent hashing is a way of spreading data across machines that keeps the reshuffling small when machines come and go.
The reason it exists is this: with ordinary hashing (hash(key) % number_of_nodes), adding one node to a ten-node cluster changes the answer for nearly every key, so nearly all your data has to move. With consistent hashing, adding that node moves roughly a tenth of it. When your cluster size changes routinely - which it does for distributed caches and databases - that difference is everything.
The trade-off: more complex than modulo hashing, and without care the distribution can be lumpy (see #36).
23. Object Storage
Object storage holds large unstructured files as objects you look up by key. It scales effectively without limit, costs very little per gigabyte, and is extremely durable. S3 is the example everyone knows.
It is the right home for images, video, documents, and backups. It is not a database, and the recurring mistake is treating it like one: there are no transactions, no rich queries, and no cheap way to read a small piece from the middle of a large file.
The trade-off: unbeatable for storing blobs, useless for asking questions about them.
24. Data Partitioning
Partitioning is the general practice of splitting data into separate parts. Sharding is one kind - splitting across machines. Splitting by time, so each month lives in its own table, is another.
You partition for three reasons: queries get faster because they scan less, old data gets easy to archive or drop a whole partition at a time, and load spreads out. The partition key decides which part each row goes into, and should be chosen around the query you run most.
The trade-off: queries that do not line up with the partition key have to touch every partition.
25. Event Sourcing
Event sourcing stores state not as the current values but as the full sequence of events that produced them. Rather than recording that an account holds £100, you record: deposited £200, withdrew £50, withdrew £50. The current balance is what you get by replaying that list.
You get a complete audit trail for free, you can rebuild any derived view by replaying from the start, and you can ask what the state was at any past moment - genuinely useful things that are painful to bolt on afterwards.
The trade-off: reading current state means replaying a lot of events, so you end up maintaining snapshots, and now you have two things to keep correct.
Section 3: Distributed Systems
What changes when the system stops being one machine.
26. Distributed System
A distributed system is one whose parts run on several machines and coordinate over a network.
It buys capability you cannot get on one machine, and it introduces a category of problem that does not exist on one machine: the network fails, parts fail while other parts keep running, and copies of data disagree. Working in distributed systems mostly means internalising three facts - the network is unreliable, clocks on different machines disagree, and anything can fail at any moment - and designing as though they are true, because they are.
The trade-off: capability and resilience, bought with a permanent increase in the number of ways things can go wrong.
27. CAP Theorem
The CAP theorem says a distributed system can guarantee at most two of: consistency, availability, and partition tolerance.
In practice that framing misleads people, because partitions - the network splitting so some machines cannot reach others - are not optional. They happen. So partition tolerance is mandatory, and the real question is narrower: when the network splits, do you choose consistency or availability?
network partition happens │ ┌─────────────┴─────────────┐ ▼ ▼ choose CONSISTENCY choose AVAILABILITY refuse the request answer anyway rather than answer with possibly stale data with stale data │ │ bank balance, like counts, inventory, seats follower counts
Most real systems decide this per piece of data rather than once for the whole system. A bank balance refuses to be wrong. A like count is happy to be briefly behind.
The trade-off: this is the trade-off. There is no configuration that gives you both during a partition.
28. Strong Consistency
Strong consistency means every read sees the most recent write, whichever machine answers it. Getting there requires the nodes to coordinate on every write, and coordination takes time.
It is what you want when stale data causes a real problem: account balances, stock levels, seats for a concert that will sell out in ninety seconds.
The trade-off: higher latency always, and reduced availability exactly when things are already going wrong.
29. Eventual Consistency
Eventual consistency means the copies will agree with each other eventually, but may disagree briefly after a write. Read immediately after writing and you might get the old value, if your read lands on a replica that has not caught up yet.
It is cheaper and stays available under conditions where strong consistency would refuse to serve. The skill is not preferring one or the other - it is being deliberate about which data can tolerate being briefly wrong and which cannot.
The trade-off: availability and speed, paid for with a window where different users see different things.
30. Consensus
Consensus is the problem of getting several machines to agree on one value despite some of them failing. It underpins leader election, distributed locks, and replicated logs.
Algorithms like Raft and Paxos solve it, and they all work on the same principle: a majority must agree before anything is committed. That majority requirement is what prevents split-brain, where two machines each believe they are in charge and both accept conflicting writes.
The trade-off: several network round trips per decision, which is why you use it only where agreement genuinely must be guaranteed.
31. Leader Election
Leader election is how a group of machines agrees which one of them is in charge - handling writes, coordinating work, or owning a resource that must not be owned twice.
It needs consensus underneath it, for the reason above: two leaders accepting writes is usually worse than no leader at all.
The trade-off: when the leader dies, there is a window - seconds, typically - where a new election is running and nothing leader-dependent can happen. You cannot design that window away, only shorten it.
32. Idempotency
An operation is idempotent if doing it twice leaves you in the same state as doing it once. GET is naturally idempotent - reading something twice changes nothing. POST is not - create something twice and you have two of them.
This matters because distributed systems cannot promise a message arrives exactly once. The request can vanish, the response can vanish, or the caller can give up waiting and retry without knowing whether the first attempt worked. Idempotent operations are safe to retry. Non-idempotent operations that get retried are how one order becomes three.
The trade-off: designing for idempotency takes deliberate effort up front, and prevents a class of bug you otherwise find in production.
33. Idempotency Key
An idempotency key is a unique value the client attaches to a request so the server can spot a duplicate. The server remembers which keys it has handled; if the same key turns up again, it returns the original result rather than doing the work twice.
This is the standard way to make a non-idempotent operation safe to retry, and it is why payment APIs almost universally require one. The important detail: the key must be generated by the client, before the first attempt, and reused on every retry. A key generated per-attempt does nothing at all.
The trade-off: the server has to store processed keys somewhere, for long enough to outlive any retry.
34. Two-Phase Commit (2PC)
Two-phase commit makes one transaction atomic across several databases or services. First the coordinator asks every participant "can you commit?". If they all say yes, it tells them all to commit. If any says no, everyone aborts.
It works, and it has a nasty failure mode: between the two phases, participants are holding locks and waiting. If the coordinator dies in that window, they wait indefinitely, still holding the locks. That fragility is why large distributed systems usually reach for sagas instead.
The trade-off: true atomicity across services, at the price of a blocking failure mode.
35. Saga Pattern
A saga is the alternative to distributed transactions. Break the operation into a series of local transactions, and give each one a compensating action that undoes it.
place order: charge payment ──▶ reserve stock ──▶ create shipment ✗ fails │ refund payment ◀── release stock ◀─────────┘ (compensate) (compensate)
Nothing is ever locked across services. If a later step fails, you run the compensations backwards. You reach consistency through undo rather than through atomicity.
The trade-off: there is a visible window where the system is half-done, and every step needs a compensating action that actually works - refunding is easy, un-sending an email is not.
36. Consistent Hashing Ring
The ring is how consistent hashing is actually built. Arrange all possible hash values in a circle, give each node responsibility for an arc of it, and store each key on the first node clockwise from where the key hashes to. Add or remove a node and only the keys in the neighbouring arc move.
the ring, unrolled (it wraps around from the right back to the left) 0 ─────────── A ──────────── B ──────────── C ─────────── 2^32 │ │ │ └─ owns keys └─ owns keys └─ owns keys up to A A → B B → C key k hashes here ─┐ ▼ 0 ─────────── A ────────·─── B ──────────── C ─────────── 2^32 └──▶ stored on B, the first node clockwise add node D between A and B: 0 ─────────── A ────── D ──── B ──────────── C ─────────── 2^32 ▲ └─ ONLY the keys between A and D move. Everything else stays exactly where it was.
Plain rings distribute unevenly, because a few random points on a circle are rarely spaced nicely. The fix is virtual nodes: each physical machine takes many small positions around the ring instead of one large one, which smooths the distribution and lets you add capacity gradually.
The trade-off: more bookkeeping, and lookups need the ring state, which every client must have.
37. Clock Skew
Clock skew is the difference between the clocks on two machines. Even with time synchronisation running, they drift, and can disagree by milliseconds or considerably more.
This breaks anything that uses timestamps from different machines to decide what happened first. Two events whose timestamps are five milliseconds apart, recorded on different machines, may genuinely have happened in the other order. You cannot fix this with better synchronisation - only reduce it.
The trade-off: you have to stop using wall-clock time for ordering and use something that tracks causality instead.
38. Vector Clock
A vector clock tracks the causal order of events across machines without trusting any clock. Every node keeps a counter; the vector holds each node's counter as that node last knew them.
When a node sends a message it includes its vector; the receiver updates its own by taking the highest value in each position. What you get out of this is not when things happened but whether one event could possibly have caused another - which is exactly the information you need to resolve conflicting writes in an eventually consistent store.
The trade-off: the vector grows with the number of nodes, and it tells you two events were concurrent without telling you which one to keep. That decision is still yours.
Section 4: Messaging and Communication
How parts of a system talk to each other when they are not in the same process.
39. Message Queue
A queue holds messages from producers until consumers are ready for them. The producer sends and moves on without waiting; the consumer works through the backlog at whatever pace it can manage.
This decouples the two sides in time, which is the valuable part. A traffic spike goes into the queue instead of into your database. The two sides can be scaled independently, and a consumer being down means a growing queue rather than lost requests.
The trade-off: the work happens later, so there is no result to return to the user immediately. Anything that needs an answer now cannot go through a queue.
40. Pub/Sub
Publish-subscribe is the other messaging shape. Producers publish to a topic, and every subscriber gets its own copy. A queue delivers each message to one consumer; pub/sub fans it out to all of them.
┌──▶ inventory service order placed ──▶ topic ───┼──▶ notification service └──▶ analytics service (the order service does not know any of these exist)
It is what you want when several unrelated parts of the system need to react to the same thing happening. The order service publishes "order placed" and carries on, with no knowledge of who is listening - which means you can add a fourth listener without touching it.
The trade-off: nobody knows who is listening, which is wonderful for decoupling and unpleasant when you are trying to work out what happens after an event.
41. Dead-Letter Queue
A dead-letter queue is where messages go after failing to process too many times. Rather than a single bad message being retried forever and blocking everything behind it, the queue sets it aside once the retry limit is hit.
Without one, a single malformed message can stall an entire pipeline. With one, the bad message is parked for a human to look at and the rest keeps moving.
The trade-off: the dead-letter queue only helps if somebody is watching it. An unmonitored one is a silent data-loss mechanism.
42. Backpressure
Backpressure is a consumer's ability to tell a producer to slow down when it cannot keep up.
Without it, a fast producer and a slow consumer produce an ever-growing queue, and the system runs until memory is exhausted and then dies - having given no warning, because from the producer's point of view everything was fine right up to the end. Backpressure makes overload visible and survivable instead. It is usually implemented with bounded queues that refuse new messages when full, which forces the producer to wait.
The trade-off: the producer now has to handle being told no, which means real error handling on a path that used to be fire-and-forget.
43. WebSocket
A WebSocket is a connection that stays open, in both directions, between client and server. Once it is established either side can send at any time, without paying to set up a new connection per message.
It is the right tool for genuinely interactive things: chat, collaborative editing, multiplayer games, live trading.
The trade-off: every connection is long-lived and stateful, which is awkward at scale. Scaling horizontally means keeping a registry of which server holds which connection, plus a pub/sub backplane to get a message to the right one - a real amount of machinery.
44. Server-Sent Events (SSE)
SSE streams updates one way, server to client, over a single HTTP connection held open. The server can push whenever it likes; the client cannot send anything back down the same channel.
It is markedly simpler than WebSockets and reconnects automatically, which makes it the better choice whenever only the server needs to push: notifications, live prices, dashboards, and streaming an AI model's response token by token.
The trade-off: one-way only - which for these use cases is not a limitation but the reason it is simpler.
Section 5: Reliability, Performance, and the Modern Set
How systems stay up when they are under strain, and the AI vocabulary that has become standard.
45. Circuit Breaker
A circuit breaker watches the calls you make to some dependency and stops making them once the failure rate crosses a threshold. It has three states:
CLOSED ──── failures exceed threshold ────▶ OPEN ▲ │ │ wait a while │ │ └──── trial calls succeed ──── HALF-OPEN ◀─┘ │ still failing └──────▶ OPEN
Closed is normal. Open means the dependency is considered down and calls are rejected immediately without even trying. Half-open lets a few trial calls through to see whether it has recovered.
The reason it exists: without one, a slow database causes your application servers to pile up waiting requests until they run out of threads or memory, and now the application is down too - and worse, the pile of retries is keeping the database from recovering. With one, the application fails fast, keeps its resources, and gives the database space to come back.
The trade-off: while the breaker is open you are deliberately failing requests that might have worked.
46. Rate Limiting
Rate limiting caps how many requests a client can make in a window, rejecting the rest until it resets. It protects you from overload and abuse, and keeps one heavy user from crowding out everyone else.
The usual implementation is the token bucket: each client has a bucket that refills at a steady rate and drains as they make requests. It allows a short burst - useful, because real traffic is bursty - while still enforcing an average over time.
The trade-off: set it too low and you break legitimate users; too high and it does not protect you. There is no correct value you can derive, only one you tune.
47. Load Shedding
Load shedding is deliberately refusing some requests when you are overloaded. Rather than attempting everything, doing all of it badly, and collapsing, you serve what you can serve properly and turn the rest away with a clear "try again shortly".
The principle behind it is that partial availability beats total failure. A system successfully serving 80% of requests is far more useful than one that attempts 100% and fails all of them. Shedding should be prioritised - drop the low-value work first and keep capacity for whatever actually matters, which usually means checkout and login rather than recommendations.
The trade-off: you are choosing to fail some users. Doing it well means deciding in advance which ones.
48. Bloom Filter
A bloom filter is a very small data structure that answers one question: is this item in the set? It gives you a definite no or a probable yes, with a false-positive rate you can tune by spending more memory.
That one-sided uncertainty is exactly what makes it useful. A database checks a bloom filter before going to disk for a key: a definite no means skip the disk read entirely, and a maybe just means you do the read you would have done anyway. Web crawlers use them to track visited URLs, where re-crawling a page occasionally is cheap and storing every URL is not.
The trade-off: you cannot remove items from a standard bloom filter, and you cannot ever be certain of a yes.
49. Embedding
An embedding is a piece of data - usually text, sometimes an image - turned into a list of numbers that represents its meaning. Things that mean similar things end up with lists that are numerically close together.
That property is what makes modern search work. Because closeness in this space tracks closeness in meaning, a system can find documents relevant to a question even when they share not a single word with it. Semantic search, recommendations, and the retrieval step in RAG are all this one idea.
The trade-off: the numbers only capture what the model that produced them understood. Change the model and every stored embedding has to be regenerated, because vectors from two different models are not comparable.
50. RAG (Retrieval-Augmented Generation)
RAG improves a language model's answers by fetching relevant information at question time and handing it to the model as context - rather than relying entirely on what the model absorbed during training.
There are two phases:
INGESTION (once, ahead of time) documents ─▶ split into chunks ─▶ embed each chunk ─▶ vector database QUERY TIME (per question) question ─▶ embed it ─▶ search vector db ─▶ top matching chunks │ ┌───────────────────────────┘ ▼ prompt = question + retrieved chunks ─▶ model ─▶ answer
The point is that the answer is grounded in documents you supplied and can check, rather than produced from general training knowledge and hoped to be right.
The trade-off: the answer is only as good as what retrieval found. Retrieve the wrong chunks and the model will confidently build an answer on them, which is a harder failure to spot than a model that simply did not know.
Key Takeaways
Grouping the fifty by what they are actually for:
| Group | Concepts | What they collectively answer |
|---|---|---|
| Infrastructure | scalability, vertical/horizontal scaling, load balancer, latency, throughput, CDN, DNS, API gateway, reverse proxy | How traffic reaches the system and how it grows |
| Data and storage | databases, ACID, indexes, sharding, replication, caching patterns, consistent hashing, object storage, partitioning, event sourcing | Where data lives and how it stays correct at scale |
| Distributed systems | CAP, consistency models, consensus, leader election, idempotency, 2PC, sagas, clock skew, vector clocks | What breaks once there is more than one machine |
| Messaging | queues, pub/sub, dead-letter queues, backpressure, WebSockets, SSE | How components talk without blocking each other |
| Reliability and modern | circuit breaker, rate limiting, load shedding, bloom filters, embeddings, RAG | How systems survive stress, and how AI fits in |
A few things worth carrying away:
No concept is free, and none exists alone. Every one on this list solves a problem by creating a smaller one somewhere else. Caching trades freshness for speed. Sharding trades query flexibility for write capacity. Replication trades correctness-right-now for read capacity. Knowing what a thing costs is what separates using it deliberately from using it because it was in a blog post.
The interview question is usually the trade-off. "What is a cache" is the warm-up. "What happens when the cached value is stale and someone has just written to it" is the actual question.
The vocabulary has shifted. Embeddings and RAG are now standard system design vocabulary, about as common in interviews as caching and sharding were a few years ago. That is a genuine change in what the field expects you to know, not a passing fashion.
System design has a language, and fluency in it is most of the difference between engineers who take part in design conversations and engineers who nod and hope. These fifty words cover the large majority of what comes up in reviews and interviews. Learn what each one does, why it exists, and what it costs, and the conversations that used to be impenetrable turn into ones where you have something to say.


