System Design
System design is the craft of turning vague product requirements into an architecture that is fast, reliable and affordable at scale, and explaining the trade-offs you made. This page covers the fundamentals, the building blocks, a repeatable interview framework, blueprints for the classic questions, a mobile and embedded angle, and what senior and principal interviewers expect.
- Never start by drawing boxes: clarify functional and non-functional requirements, then estimate scale.
- Scale out with stateless services behind load balancers; push state into caches, databases and queues.
- The universal levers: cache, replicate for reads, shard for writes, go async with queues, put static content on a CDN, remove single points of failure.
- Under a network partition you choose consistency or availability (CAP); even without one you trade latency for consistency (PACELC).
- Make retries safe with idempotency keys, and protect services with timeouts, backoff with jitter, rate limits and circuit breakers.
- Senior candidates drive the conversation, quantify, name every trade-off, and cover failure modes, observability and rollout.
Fundamentals: scalability, latency and availability
Every design discussion comes back to a few mental models: how the system grows (scalability), how fast it answers (latency) and how much work it can do (throughput), and how often it is up (availability).
A restaurant getting popular. Vertical scaling is hiring one super-chef who cooks faster, until no chef can go any faster and the restaurant closes whenever that chef is sick. Horizontal scaling is opening more kitchens with ordinary chefs and a host at the door sending diners to whichever kitchen is free. Latency is how long one diner waits for a meal; throughput is how many meals the restaurant serves per hour. Availability is the fraction of the year the doors are open. The host is the load balancer, the kitchens are stateless servers, and a spare kitchen is redundancy that removes a single point of failure.
Vertical vs horizontal scaling
Vertical (scale up)
- Bigger machine: more CPU, RAM, faster disks
- Simple: no code changes, no distribution
- Hard ceiling and rising cost per unit
- Single point of failure; upgrades often need downtime
Horizontal (scale out)
- More machines behind a load balancer
- Near-linear, elastic growth; tolerates failures
- Needs stateless services and data partitioning
- Adds network hops, coordination and operational complexity
Statelessness is the enabler: if app servers keep no session state (sessions live in Redis, a database or a signed token such as a JWT), any server can handle any request, so you can add, remove or replace servers freely and fail over instantly.
Latency vs throughput
- Latency is the time to complete one request. Throughput is requests (or bytes) per second. Batching increases throughput but can increase latency; you optimise the one the product needs.
- Report latency as percentiles, not averages: p50 (median), p95, p99, p99.9. The average hides the slow tail.
- Tail latency amplification: if one page load fans out to 100 backend calls and each has a 1 percent chance of being slow, about 63 percent of page loads hit at least one slow call (1 - 0.99^100). Mitigations: hedged or backup requests, timeouts, fewer fan-out hops, caching.
- Little's law: concurrent requests in flight = arrival rate x average latency. At 1,000 requests per second and 200 ms latency, about 200 requests are in flight, which sizes your thread pools and connection pools.
Availability and the nines
| Availability | Downtime per year | Downtime per month |
|---|---|---|
| 99% (two nines) | about 3.65 days | about 7.3 hours |
| 99.9% (three nines) | about 8.8 hours | about 44 minutes |
| 99.99% (four nines) | about 53 minutes | about 4.4 minutes |
| 99.999% (five nines) | about 5.3 minutes | about 26 seconds |
- Components in series multiply: a request that needs a 99.9% service and a 99.9% database is only about 99.8% available.
- Redundant components in parallel: two independent 99% replicas give 1 - 0.01^2 = 99.99%, provided failures are truly independent (different racks, zones or regions).
- Single point of failure (SPOF): any component whose failure takes the whole system down. Remove it with redundancy, health checks and automatic failover.
- SLI, SLO, SLA: an SLI is a measured indicator (for example, the fraction of requests under 300 ms); an SLO is the internal target (99.9% of requests under 300 ms); an SLA is the contractual promise to customers with penalties. The gap between 100% and the SLO is the error budget that teams can spend on risky launches.
- RTO and RPO: recovery time objective (how long you may be down) and recovery point objective (how much recent data you may lose) drive backup and multi-region choices.
Consistency, CAP and PACELC
As soon as data is copied to more than one machine, copies can disagree. Consistency models define what a reader is allowed to see; CAP and PACELC describe what you give up to get stronger guarantees.
Two bank branches that share one ledger by phone. If the phone line is cut (a network partition), each clerk must choose: refuse to process withdrawals until the line is back (consistency) or keep serving customers and reconcile the books later (availability). Even when the line works, calling the other branch before every transaction makes customers wait longer (the latency cost in PACELC). The clerks are replicas, the phone line is the network, and the ledger is the replicated data.
CAP theorem
In a distributed system, when a network Partition happens you can keep only one of Consistency (every read sees the latest write, specifically linearizability) or Availability (every request to a non-failed node gets a non-error response). Because partitions cannot be ruled out in real networks, the practical choice is CP or AP during a partition.
- CP: reject or block requests on the minority side rather than diverge. Examples: bank ledgers, inventory reservation, leader election and configuration stores such as ZooKeeper or etcd.
- AP: keep serving and reconcile later (eventual consistency, conflict resolution). Examples: shopping carts, social feeds, like counts, DNS, Cassandra or DynamoDB in their default modes.
- "CA" only exists when there is no partition, which in practice means a single node or a single tightly coupled site.
PACELC
PACELC extends CAP: if there is a Partition, choose A or C; Else (normal operation), choose Latency or Consistency. Synchronous replication across regions is consistent but slow; asynchronous replication is fast but readers may see stale data. DynamoDB and Cassandra are typically PA/EL; Spanner and most relational setups with synchronous replicas are PC/EC.
Consistency models (strongest to weakest)
| Model | Guarantee | Typical use |
|---|---|---|
| Linearizable (strong) | Every operation appears to take effect at one instant between its start and end; everyone sees the same latest value | Locks, leader election, balances, unique usernames |
| Sequential | All nodes see operations in the same order, consistent with each client's own order, but not necessarily real time | Some replicated logs |
| Causal | Operations that are causally related (a reply after a post) are seen in order by everyone; unrelated ones may differ | Comments and replies, collaborative apps |
| Read-your-writes | A client always sees its own writes | Profile edits, settings |
| Monotonic reads | A client never sees data go backwards in time | Pinning a user to one replica |
| Eventual | If writes stop, all replicas converge eventually | Likes, view counts, feeds, DNS |
Quorums
With N replicas, a write waits for W acknowledgements and a read queries R replicas. If R + W > N, every read set overlaps every write set, so a read sees at least one copy of the latest write. N = 3, W = 2, R = 2 is a common balance; W = 1, R = 1 is fast but only eventually consistent; W = N makes writes fail if any replica is down. Conflicts in leaderless systems are resolved with version vectors, last-writer-wins timestamps (can lose updates) or CRDTs (data types that merge automatically).
ACID vs BASE
ACID (typical of relational databases)
- Atomicity: all or nothing
- Consistency: constraints always hold
- Isolation: concurrent transactions do not interfere (read committed, repeatable read, serializable)
- Durability: committed data survives crashes (write-ahead log)
BASE (typical of large NoSQL stores)
- Basically available: stays up under failures
- Soft state: state can change without new input as replicas converge
- Eventual consistency
- Trades guarantees for scale and availability
Isolation anomalies to know: dirty read, non-repeatable read, phantom read, lost update and write skew. Serializable isolation prevents all of them at a performance cost; many databases default to read committed or snapshot isolation.
Back-of-envelope estimation
Estimation turns requirements into numbers that justify your design: is it read-heavy or write-heavy, does the data fit on one machine, how many servers do you need? Precision does not matter; stating assumptions and rounding aggressively does.
Planning a wedding. You do not need the exact guest count to book a venue; "about 200 guests, two drinks each, 10 percent no-shows" is enough to choose a hall, order catering and hire staff. In system design the guests are daily active users, the drinks are requests per user, and the hall size is your storage and server count. Getting within a factor of two is good enough to pick the right architecture.
Numbers to memorise
| Quantity | Value |
|---|---|
| Seconds per day | 86,400 (round to 10^5) |
| 1 million requests per day | about 12 per second on average |
| Peak traffic | about 2-3x average (more for spiky events) |
| Powers of ten | thousand 10^3 (KB), million 10^6 (MB), billion 10^9 (GB), trillion 10^12 (TB) |
| Common sizes | int 4 bytes, long or timestamp 8 bytes, UUID 16 bytes, ASCII char 1 byte, short text post about 300 bytes, image about 200 KB-2 MB, minute of HD video about 50-100 MB |
| One commodity server | roughly thousands to tens of thousands of simple requests per second; tens to hundreds of GB of RAM |
Latency numbers every engineer should know (order of magnitude)
| Operation | Approximate time |
|---|---|
| L1 cache reference | about 1 ns |
| Main memory reference | about 100 ns |
| Read 1 MB sequentially from memory | about 10-50 microseconds (classic 2012 number was ~250 μs) |
| SSD random read | about 16-100 microseconds |
| Round trip within one data centre | about 0.5 ms |
| Read 1 MB sequentially from SSD | about 0.2-1 ms |
| Hard disk seek | about 10 ms |
| Send 1 MB over a 1 Gbps link | about 10 ms |
| Round trip California to Netherlands | about 150 ms |
Takeaway: memory is far faster than SSD, SSD far faster than disk, and anything crossing a region costs tens to hundreds of milliseconds. Keep hot data in memory and keep cross-region round trips off the request path.
Formulas
- QPS = daily active users x actions per user per day / 86,400. Peak QPS = QPS x peak factor.
- Storage = new items per day x bytes per item x retention days x replication factor.
- Bandwidth = QPS x average response size (separately for ingress and egress).
- Cache size (80/20 rule) = about 20 percent of the daily read working set.
- Servers = peak QPS / QPS per server, plus headroom (for example 30 percent) and N+1 or N+2 for failures.
Worked example: URL shortener
Assume: 100 M new short URLs per day, read:write = 10:1, keep 5 years, ~500 bytes per record
Writes: 100 M / 86,400 ~ 1,160 per second (peak ~ 3,000)
Reads: 10 x 1,160 ~ 11,600 per second (peak ~ 30,000)
Storage: 100 M x 500 B x 365 x 5 ~ 91 TB (x3 replication ~ 275 TB)
Keys needed: 100 M x 365 x 5 ~ 183 billion
Base62 length 7: 62^7 ~ 3.5 trillion -> 7 characters is plenty
Cache: hot-set proxy = 20% of daily *new* URLs x 500 B ~ 10 GB -> fits in RAM
(distinct URLs among 1 B daily reads is unknown; 20% of new keys is the usual interview stand-in)
Conclusions the numbers force: the system is read-heavy (so cache aggressively and use a CDN), storage outgrows one machine within a few years (so plan for a sharded key-value store), and 7-character keys are enough.
Load balancers, API gateways, DNS and CDNs
Before a request reaches your code it passes through DNS (which finds an address), often a CDN (which may answer from a nearby edge), and a load balancer or API gateway (which picks a healthy server and applies cross-cutting policy).
A large hospital. DNS is the directory that tells you which hospital building to go to. The CDN is the neighbourhood pharmacy that stocks the most common medicines so you do not have to travel to the main hospital. The load balancer is the reception desk that sends each patient to a free doctor and stops sending patients to a doctor who has gone home (health checks). The API gateway is reception plus security: it checks your ID (authentication), limits how many times you can visit per day (rate limiting) and sends you to the right department (routing).
Client | DNS lookup (GeoDNS / latency-based routing picks nearest region) v CDN edge --cache hit--> response (images, JS, video segments) | miss / dynamic v L4/L7 load balancer (TLS termination, health checks) v API gateway (auth, rate limit, routing, request IDs) v Stateless service instances (autoscaled) v Cache -> Database -> Queue -> Workers
Load balancers
| Aspect | Layer 4 (transport) | Layer 7 (application) |
|---|---|---|
| Decides using | IP addresses and ports | HTTP method, path, headers, cookies |
| Speed | Very fast, protocol agnostic | Slightly more CPU (parses requests) |
| Features | Connection-level balancing | Path-based routing, TLS termination, sticky sessions, retries, compression, header rewriting |
| Examples | AWS NLB, LVS | AWS ALB, Nginx, Envoy, HAProxy (HTTP mode) |
Algorithms: round robin (equal servers), weighted round robin (mixed capacity), least connections (long-lived or uneven requests), least response time, IP or consistent hashing (session affinity or cache locality), and "power of two random choices" (pick two at random, send to the less loaded, which is simple and close to optimal).
Health checks remove failing instances. The load balancer itself must not be a SPOF: run it as an active-passive or active-active pair with a floating IP, or use a managed, regionally redundant service. DNS round robin or anycast spreads traffic across load balancers or regions.
API gateway and reverse proxy
A reverse proxy sits in front of servers (Nginx). An API gateway is a reverse proxy with API-level concerns: authentication, rate limiting, request validation, routing to microservices, protocol translation (REST to gRPC), request aggregation, and adding correlation IDs for tracing. A service mesh (Envoy sidecars with Istio or Linkerd) moves similar concerns (mTLS, retries, timeouts, metrics) to service-to-service traffic.
DNS
DNS resolves names to IPs through a hierarchy (resolver, root, top-level domain, authoritative server) with caching at every level controlled by TTLs. GeoDNS or latency-based routing sends users to the nearest region; low TTLs allow faster failover but increase DNS traffic. Anycast announces the same IP from many locations so routing delivers users to the nearest one.
CDN
A content delivery network caches content on edge servers close to users, cutting latency and origin load. It is essential for images, video, JavaScript and CSS at global scale, and can also cache some API responses, terminate TLS and absorb DDoS traffic.
- Pull CDN: the edge fetches from the origin on the first miss and caches by TTL. Simple; first request is slow.
- Push CDN: you upload content to the CDN ahead of time. Good for large, predictable assets such as video releases.
- Invalidation: use versioned (content-hashed) file names such as
app.3f9a2c.jsso you never need to purge; purge APIs are slow and eventually consistent. - Signed URLs protect private content on the CDN with expiring tokens.
Caching
A cache stores the results of expensive work (database queries, computed pages, remote calls) in fast storage, usually memory, so later requests are served quickly. Caching appears at every layer: browser, CDN, API gateway, application (in-process), distributed cache (Redis, Memcached), and inside the database (buffer pool).
Keeping your most-used tools on your desk instead of walking to the warehouse each time. The desk is small (limited memory), so you must decide which tools to put back when it is full (eviction policy). If a colleague replaces a tool in the warehouse, the copy on your desk is now outdated (stale data), so you either check the date on it (TTL) or ask to be told when it changes (invalidation). If everyone's desk copy disappears at the same moment, the whole office stampedes to the warehouse at once (cache stampede).
Read and write strategies
| Strategy | How it works | Pros | Cons |
|---|---|---|---|
| Cache-aside (lazy loading) | App reads cache; on miss, reads DB and fills cache | Most common; cache holds only requested data; cache failure is not fatal | First read is slow; stale until TTL or invalidation |
| Read-through | Cache library loads from DB on miss itself | Simpler app code | Same staleness issues; cache must know the DB |
| Write-through | Write to cache and DB synchronously | Cache always fresh | Slower writes; caches data that may never be read |
| Write-back (write-behind) | Write to cache; flush to DB asynchronously | Very fast writes; batching | Data loss if cache node dies before flush |
| Write-around | Write to DB only; cache filled on later reads | Avoids polluting cache with write-once data | Recent writes miss the cache |
| Refresh-ahead | Refresh hot keys before they expire | Avoids miss latency for hot data | Wasted work if predictions are wrong |
# cache-aside read and invalidate-on-write
def get_user(user_id):
key = f"user:{user_id}"
user = cache.get(key)
if user is None:
user = db.query_user(user_id)
cache.set(key, user, ttl=300 + random.randint(0, 60)) # jitter
return user
def update_user(user_id, fields):
db.update_user(user_id, fields) # 1. write the source of truth
cache.delete(f"user:{user_id}") # 2. delete, don't set, the cached copy
Deleting rather than updating the cache avoids a race where two concurrent writers leave the older value in the cache. A small race still exists (a reader loads the old value just before the delete completes); short TTLs, delayed double-delete, or streaming invalidations from the database change log (change data capture) narrow it further.
Eviction policies
- LRU (least recently used): evict the item not used for the longest time. Good default for temporal locality; O(1) with a hash map plus doubly linked list.
- LFU (least frequently used): evict the least-used item. Better for stable popularity; needs frequency decay so old favourites do not stay forever.
- FIFO and random: simple, sometimes good enough.
- TTL: expire after a fixed time, which bounds staleness regardless of eviction.
Cache problems and fixes
| Problem | What happens | Fixes |
|---|---|---|
| Stampede (thundering herd) | A hot key expires and thousands of requests hit the DB at once | Request coalescing or a per-key lock (only one rebuilds), stale-while-revalidate, TTL jitter, early probabilistic refresh |
| Penetration | Requests for keys that do not exist always miss and hit the DB | Cache negative results briefly; Bloom filter of existing keys |
| Avalanche | Many keys expire together, or the cache cluster restarts cold | Randomised TTLs, warm-up, replicas, circuit breaker to protect the DB |
| Hot key | One key (a celebrity profile) overloads one cache shard | Replicate the key across shards with suffixes, local in-process cache, CDN |
| Staleness and inconsistency | Cache and DB disagree | TTL, invalidate on write, versioned keys, change data capture |
Redis
- Rich data types: strings, hashes, lists, sets, sorted sets, streams
- Persistence options, replication, Redis Cluster sharding
- Atomic operations and Lua scripts (rate limiters, leaderboards, locks)
- Mostly single-threaded command execution
Memcached
- Simple key-value strings only
- Multi-threaded, very fast, easy to scale out
- No persistence or replication built in
- Good for plain look-aside caching
Databases, indexes, replication and sharding
The database is usually the hardest part to scale because it holds state. Choose the storage model from the access patterns, then scale reads with replication and caching, and scale writes and storage with sharding.
A national library. An index is the card catalogue: instead of walking every shelf, you look up the title and go straight to the right shelf, but every new book also needs a new catalogue card (slower writes). Replication is opening branch libraries with copies of popular books so more people can read at once (read scaling), though a new edition reaches the branches a little late (replication lag). Sharding is splitting the collection across buildings by subject or author surname so no single building has to hold everything (write and storage scaling); a question spanning two buildings (a cross-shard join) means a trip to both.
SQL vs NoSQL
| Type | Examples | Strengths | Use for |
|---|---|---|---|
| Relational (SQL) | PostgreSQL, MySQL | Schemas, joins, ACID transactions, mature tooling | Payments, orders, users, anything relational with integrity constraints |
| Key-value | DynamoDB, Redis, RocksDB | Simple, very fast lookups by key; easy horizontal scale | Sessions, URL mappings, carts, feature flags |
| Document | MongoDB, Couchbase | Flexible JSON-like records, nested data | Catalogues, content, user profiles |
| Wide-column | Cassandra, HBase, Bigtable | Massive write throughput, partition key plus clustering order | Time series, messages, activity logs |
| Graph | Neo4j, Neptune | Fast multi-hop relationship queries | Social graphs, fraud rings, recommendations |
| Search | Elasticsearch, OpenSearch | Inverted index, full-text and faceted search | Search boxes, log analytics |
| Time-series | InfluxDB, Prometheus, TimescaleDB | Compression and downsampling by time | Metrics, telemetry |
| Blob / object storage | S3, GCS | Cheap, durable storage of large files | Images, video, backups, logs |
| Columnar warehouse | BigQuery, Redshift, ClickHouse | Fast scans and aggregations (OLAP) | Analytics, dashboards |
OLTP vs OLAP: OLTP systems handle many small, indexed reads and writes with low latency (the app database). OLAP systems handle few, huge scans and aggregations (the warehouse). Keep them separate, feeding the warehouse through ETL or change data capture.
Indexes and storage engines
- A B-tree / B+ tree index keeps keys sorted in wide, shallow pages, giving O(log n) lookups and efficient range scans. Index columns used in WHERE, JOIN and ORDER BY clauses.
- Costs: every insert or update must maintain each index (slower writes), plus extra storage. Do not index everything on write-heavy tables.
- Composite indexes follow the leftmost-prefix rule: an index on (country, city) helps queries on country or on country + city, not on city alone.
- A covering index includes every column a query needs, so the table itself is never read.
- LSM trees (Cassandra, RocksDB, LevelDB) buffer writes in memory (memtable) and flush sorted files (SSTables) that are merged in the background (compaction). Writes are sequential and very fast; reads may check several files, helped by Bloom filters.
- An inverted index maps each word to the documents containing it; it powers full-text search.
Replication
| Topology | How it works | Trade-offs |
|---|---|---|
| Leader-follower (primary-replica) | All writes go to the leader, which streams changes to followers; reads can go to followers | Simple and common. Async replication means lag, so stale reads and possible loss of recent writes on failover |
| Multi-leader | Several leaders accept writes (for example one per region) and replicate to each other | Local writes in each region, but write conflicts must be resolved |
| Leaderless (Dynamo style) | Clients write to and read from several replicas using quorums | High availability; needs read repair, hinted handoff, and conflict resolution |
- Synchronous replication waits for the follower before acknowledging (durable, slower); asynchronous does not (fast, may lose data). Semi-synchronous waits for one follower.
- Read-your-writes on top of async replicas: route a user's reads to the leader for a short time after they write, or track the replication position they need.
- Failover must avoid split brain (two leaders): use a consensus-based coordinator (Raft, Paxos, ZooKeeper or etcd) and fencing tokens.
Sharding (partitioning)
Sharding splits data across nodes by a shard key so each node stores and serves a subset. It scales writes and storage beyond one machine.
| Scheme | How | Pros | Cons |
|---|---|---|---|
| Range | Key ranges per shard (A-F, G-M, ...) | Efficient range scans | Hot spots on sequential keys such as timestamps |
| Hash | hash(key) mod N or consistent hashing | Even spread | Range queries hit all shards; mod N remaps nearly everything when N changes |
| Directory / lookup | A service maps key to shard | Flexible, easy to move tenants | Lookup service is an extra dependency |
| Geographic | By user region | Data locality, residency compliance | Uneven regions; cross-region users |
Choosing a shard key: it should have high cardinality, spread load evenly, and match the dominant query so most requests touch one shard (for example user_id for a user's data, conversation_id for chat messages). Hard parts: hot shards (a celebrity), cross-shard joins and transactions (avoid them or use sagas), secondary indexes (local per shard or global), and resharding without downtime.
Consistent hashing
Map both nodes and keys onto a hash ring; each key belongs to the first node clockwise. Adding or removing a node moves only the keys between it and its neighbour, about 1/N of the data, instead of nearly all keys as with hash mod N. Virtual nodes (each physical node appears at many points on the ring) smooth out the load and let bigger machines take more points. Used by DynamoDB, Cassandra, distributed caches and load balancers.
Hash ring (positions 0 .. 2^32), walking clockwise: 0 ---- k1 ---- [N1] ---- k2 ---- k3 ---- [N2] ---- k4 ---- [N3] ---- (wraps to 0) k1 -> N1 k2, k3 -> N2 k4 -> N3 Add node N4 between k2 and k3: 0 ---- k1 ---- [N1] ---- k2 ---- [N4] ---- k3 ---- [N2] ---- k4 ---- [N3] only k2 moves (N2 -> N4); every other key stays where it was
Message queues, streams and idempotency
A message queue or log sits between producers and consumers so they no longer have to be up, fast, or scaled at the same moment. This is the backbone of asynchronous and event-driven architectures.
The ticket rail in a restaurant kitchen. Waiters (producers) clip orders to the rail and go back to their tables; cooks (consumers) take orders at their own pace. A dinner rush just makes the rail longer instead of making waiters wait at the kitchen door (buffering and backpressure). If a cook drops an order, it goes back on the rail (redelivery), so the kitchen must be able to recognise "table 12's second steak" when it sees it twice (idempotency). An order nobody can cook goes on a separate spike for the manager (dead-letter queue).
Why use one
- Decoupling: producers do not need to know who consumes or whether they are up.
- Load levelling: absorb spikes and let consumers work at a steady rate.
- Asynchrony: return to the user quickly and do slow work (emails, transcoding, fan-out) in the background.
- Fan-out: one event, many independent consumers (analytics, search indexing, notifications).
- Retries and durability: failed work is redelivered rather than lost.
Queue (RabbitMQ, SQS)
- Each message goes to one consumer and is deleted after acknowledgement
- Per-message routing, priorities, delays
- Great for task and job distribution
- Replay is not the normal model
Log / stream (Kafka, Kinesis, Pulsar)
- Append-only, partitioned, retained for days or longer
- Many consumer groups read independently at their own offsets
- Ordering guaranteed within a partition
- Replay history, event sourcing, stream processing
Delivery semantics
| Semantic | Meaning | How |
|---|---|---|
| At-most-once | May lose messages, never duplicates | Acknowledge before processing |
| At-least-once | Never loses, may duplicate | Acknowledge after processing; the common default |
| Exactly-once (effectively) | Each message affects state once | At-least-once delivery plus idempotent processing or transactional writes (for example Kafka transactions within Kafka) |
End-to-end exactly-once delivery across arbitrary systems is not achievable in general; what you build is at-least-once delivery with idempotent consumers.
Idempotency
An operation is idempotent if doing it twice has the same effect as doing it once. SET balance = 100 is idempotent; balance += 10 is not. Retries, timeouts and at-least-once queues all cause duplicates, so every write path that can be retried must be idempotent.
# idempotency key pattern for a payment API
def charge(request):
key = request.headers["Idempotency-Key"] # client-generated UUID
existing = store.get(key)
if existing:
return existing.response # replay, do not re-charge
with db.transaction():
if not store.insert_if_absent(key, status="in_progress"):
return conflict_or_wait() # a concurrent duplicate
result = payments.charge(request.amount)
store.update(key, status="done", response=result)
return result
Other techniques: unique constraints on a natural key (order ID), conditional writes with version numbers (optimistic concurrency), and consumer-side "already processed message IDs" tables.
Patterns around queues
- Dead-letter queue (DLQ): after N failed attempts, park the message for inspection instead of blocking the queue forever (a "poison message").
- Retry with exponential backoff and jitter: wait 1 s, 2 s, 4 s ... plus randomness so retries do not synchronise.
- Ordering: Kafka orders only within a partition, so partition by the entity whose order matters (for example
user_idororder_id). - Transactional outbox: write the business row and an "outbox" event row in the same database transaction; a relay publishes outbox rows to the queue. This avoids the dual-write problem where the DB commit succeeds but the publish fails, or vice versa.
- Saga: a long business transaction across services as a sequence of local transactions, each with a compensating action (refund, release inventory) if a later step fails. Choreographed through events or orchestrated by a coordinator.
- Backpressure: when consumers fall behind, slow producers (bounded queues, rate limits) or scale consumers; watch consumer lag as a key metric.
- Event sourcing and CQRS: store the log of events as the source of truth and build read-optimised views from it; separate the write model from read models.
API design and rate limiting
The API is the contract between clients and your system. Defining it early pins down what the system must do. Rate limiting protects that API from abuse, runaway clients and overload.
A restaurant menu and a nightclub bouncer. The menu (API) tells customers exactly what they can order and in what form, and once printed it is hard to change without confusing regulars (versioning and backwards compatibility). The bouncer (rate limiter) lets people in at a controlled pace: a token bucket is a bouncer who allows a small group through at once if the club has been quiet, but never more than the club's capacity over time. Returning HTTP 429 is the bouncer saying "come back in ten minutes".
API styles
| Style | Strengths | Weaknesses | Typical use |
|---|---|---|---|
| REST over HTTP/JSON | Simple, cacheable, universal tooling | Over- or under-fetching; loose contracts | Public APIs, CRUD resources |
| gRPC (HTTP/2 + Protocol Buffers) | Fast binary encoding, strict schemas, streaming, generated clients | Harder to use from browsers; less human-readable | Internal service-to-service calls |
| GraphQL | Client picks exactly the fields it needs; one round trip | Caching and cost control are harder; N+1 query risk | Mobile and web clients aggregating many resources |
| WebSocket / server-sent events | Server push, low-latency bidirectional updates | Stateful connections to manage and scale | Chat, live dashboards, presence |
| Webhooks | Server-to-server push on events | Receiver must be reachable; retries and signatures needed | Payment callbacks, integrations |
Good API design checklist
- Resource-oriented URLs and correct methods:
GET /users/{id},POST /orders,PATCHfor partial updates. GET, PUT and DELETE should be idempotent; POST is not unless you add an idempotency key. - Pagination: cursor-based (
?after=opaque_cursor&limit=50) is stable under inserts and scales; offset pagination gets slow and skips or repeats rows when data changes. - Versioning:
/v1/in the path or a header; only make additive, backwards-compatible changes within a version. - Consistent errors with meaningful status codes: 400 bad input, 401 unauthenticated, 403 forbidden, 404 not found, 409 conflict, 429 too many requests, 5xx server errors.
- Authentication (OAuth 2.0 or tokens), authorisation checks on every resource, TLS everywhere.
- Idempotency keys on retries of non-idempotent operations; request IDs for tracing.
- Batch endpoints and field filtering to cut round trips for mobile clients.
POST /v1/urls { "long_url": "...", "custom_alias": null, "expires_at": null }
-> 201 { "short_url": "https://sho.rt/aZ3kP9q" }
GET /{short_key} -> 301 or 302 Location: long_url
GET /v1/users/{id}/feed?cursor=abc&limit=20
-> 200 { "items": [...], "next_cursor": "def" }
Rate-limiting algorithms
| Algorithm | How it works | Pros | Cons |
|---|---|---|---|
| Token bucket | Bucket of capacity B refills at R tokens per second; each request takes a token | Allows short bursts up to B while enforcing average rate R; tiny state | Burst may still hurt a fragile backend |
| Leaky bucket | Requests enter a queue that drains at a fixed rate | Smooth, constant output rate | Bursts are delayed or dropped; adds latency |
| Fixed window counter | Count requests per clock window (per minute) | Simplest | Up to 2x the limit across a window boundary |
| Sliding window log | Store a timestamp per request; count those within the last window | Exact | Memory proportional to request volume |
| Sliding window counter | Weighted blend of the current and previous fixed windows | Close to exact with O(1) memory | Approximate |
class TokenBucket:
def __init__(self, rate, capacity):
self.rate, self.capacity = rate, capacity
self.tokens, self.last = capacity, time.monotonic()
def allow(self):
now = time.monotonic()
self.tokens = min(self.capacity, self.tokens + (now - self.last) * self.rate)
self.last = now
if self.tokens >= 1:
self.tokens -= 1
return True
return False # caller returns 429 with a Retry-After header
In a distributed deployment, keep the bucket state in a shared store such as Redis and update it atomically (a Lua script or an atomic increment with expiry) so concurrent gateway nodes do not race. For lower latency, keep local buckets on each node with a share of the global limit and sync periodically, trading some accuracy. Key limits by user, API key, IP or endpoint, and return 429 Too Many Requests with Retry-After and remaining-quota headers.
Reliability, resilience and observability
At scale, something is always failing: a disk, a node, a network link, a dependency, a deploy. Reliable systems assume failure and contain it so a local problem does not become a global outage.
The electrical system of a house. A circuit breaker trips when one appliance shorts, so the rest of the house keeps its power instead of the whole building catching fire. Bulkheads on a ship are watertight compartments: one flooded compartment does not sink the ship. A timeout is refusing to wait on hold forever. Retrying with backoff and jitter is calling back later at a random time rather than everyone redialling the instant the line frees up. Monitoring is the smoke detector and the electricity meter that tell you something is wrong before you smell smoke.
Resilience patterns
| Pattern | What it does | Watch out for |
|---|---|---|
| Timeouts | Bound how long you wait for any remote call | Without them, slow dependencies exhaust threads and cascade |
| Retries with exponential backoff and jitter | Recover from transient failures | Only retry idempotent operations; cap attempts; use retry budgets to avoid retry storms |
| Circuit breaker | After repeated failures, fail fast for a while, then probe (half-open) | Needs a sensible fallback |
| Bulkhead | Separate thread or connection pools per dependency | One slow dependency cannot starve the others |
| Load shedding | Reject low-priority work when overloaded | Better to serve 80 percent well than 100 percent badly |
| Graceful degradation | Serve reduced functionality (cached or default data) | Decide in advance which features are optional |
| Health checks and auto-healing | Replace unhealthy instances automatically | Distinguish liveness (restart me) from readiness (do not send traffic yet) |
| Redundancy across zones and regions | Survive data-centre failures | Active-active needs conflict handling; active-passive needs tested failover |
Distributed-systems essentials
- Consensus (Raft, Paxos) lets a group of nodes agree on a value or a leader despite failures, as long as a majority is alive; with 5 nodes you tolerate 2 failures. Used for leader election, configuration and locks (etcd, ZooKeeper, Consul).
- Clocks are unreliable: wall clocks drift and jump, so do not order events across machines by timestamp alone. Use logical clocks (Lamport), version vectors, or bounded-uncertainty clocks (Spanner's TrueTime).
- Unique IDs at scale: UUIDs (random, 128-bit, not sortable), Snowflake-style IDs (timestamp + machine ID + sequence; 64-bit and roughly time-ordered), or database ticket servers that hand out ranges.
- Distributed locks need leases with expiry and fencing tokens (monotonic numbers checked by the storage layer) so a paused former holder cannot corrupt data.
- Two-phase commit gives atomic cross-node transactions but blocks if the coordinator fails; microservices usually prefer sagas.
- Gossip protocols spread membership and failure information peer to peer.
Monolith vs microservices
Monolith
- One deployable; simple to build, test and debug
- In-process calls, easy transactions
- Best for small teams and early products
- Scales as a unit; large codebases slow teams down
Microservices
- Independent deploy and scaling per service; team autonomy
- Fault isolation and technology choice per service
- Network calls, partial failures, distributed data and transactions
- Needs strong observability, CI/CD and platform tooling
Observability
- Metrics: the RED method for services (Rate, Errors, Duration) and the USE method for resources (Utilisation, Saturation, Errors). The four golden signals: latency, traffic, errors, saturation.
- Logs: structured (JSON), with request IDs, sampled at high volume, scrubbed of personal data.
- Traces: distributed tracing (OpenTelemetry) follows one request across services to find where time goes.
- Alerts on SLO burn rate and user-facing symptoms, not on every internal cause.
- Safe rollout: feature flags, canary deployments (1 percent, then 10, then 50, then 100), automatic rollback on SLO regression, and a kill switch.
The step-by-step design framework
A system design interview is an open-ended conversation, usually 45 to 60 minutes. The biggest failure mode is not a wrong answer but a lack of structure: drawing boxes before understanding the problem, or spending 30 minutes on one detail. Use the same framework every time and say each step out loud.
An architect designing a house. First they ask how many people will live there, the budget and the style (requirements). Then they work out square footage and rooms (estimation), agree how people will move through the house (API), decide what materials to use (data model), sketch the floor plan (high-level design), detail the tricky parts like the staircase or the foundation on a slope (deep dive), and finally check fire exits, plumbing failures and future extensions (bottlenecks, failure modes and evolution). Nobody pours concrete before the floor plan is agreed.
- Requirements (5-8 min) Functional: the three to five core features, and what is out of scope. Non-functional: users and scale, read/write ratio, latency target (p99), availability, consistency needs, durability, regions, privacy and compliance. Ask; do not assume.
- Estimation (3-5 min) Daily active users, read and write QPS (average and peak), storage per year, bandwidth, cache size. Each number should drive a decision.
- API design (3-5 min) A few core endpoints or RPCs with request and response shapes. This pins down the contract and the data you must store.
- Data model (3-5 min) Entities, relationships, access patterns, SQL vs NoSQL with a reason, primary keys and partition keys.
- High-level design (8-10 min) Draw clients, DNS/CDN, load balancer, services, cache, databases, queues and workers. Walk through the main read and write flows end to end.
- Deep dive (10-15 min) Pick one or two of the hardest parts, or follow the interviewer's lead: sharding, the feed fan-out, consistency, the hot path, the ID generator.
- Bottlenecks and trade-offs (5 min) Single points of failure, hot keys, scaling limits, what breaks at 10x. Name each trade-off: consistency vs availability, latency vs cost, complexity vs flexibility.
- Operability and rollout (2-3 min) Metrics, alerts, dashboards, canary rollout, feature flags, kill switch, cost, and how the design evolves.
Generic high-level design
+-------+
Clients --DNS--> | CDN | (static, media)
| +-------+
v
+---------------+ +--------------------+ +------------------+
| Load balancer | --> | Stateless services | --> | Cache (Redis) |
+---------------+ | (autoscaled) | +------------------+
+--------------------+ --> | DB primary |--> read replicas
| +------------------+
v (sharded by key)
+----------------+ +---------+ +----------------+
| Queue / stream | --> | Workers | --> | Blob store, |
+----------------+ +---------+ | search, OLAP |
+----------------+
Universal levers to pull
Cache
Read-heavy and tolerant of slight staleness? Add a cache layer and a CDN.
Replicate
Too many reads for one database? Add read replicas; handle replication lag.
Shard
Too many writes or too much data? Partition by a well-chosen key.
Go async
Slow or bursty work? Put it on a queue and process it with workers.
Denormalise
Expensive joins on the hot path? Precompute views, such as timelines.
Remove SPOFs
Redundancy across zones, health checks and automatic failover.
Blueprints I: infrastructure designs
These blueprints are compact starting points. In an interview, run each through the full framework; here the focus is on the core idea, the crux, and what to deep-dive.
Blueprints are like chess openings. Strong players memorise the first moves of common openings so they reach a sensible position quickly, then think hard about the middle game. Knowing that a URL shortener is "key generation plus a read-heavy key-value store" gets you to the interesting trade-offs (collisions, caching, analytics) in five minutes instead of twenty. The opening is the blueprint; the middle game is the deep dive on the interviewer's chosen weak spot.
URL shortener (TinyURL, bit.ly)
- Requirements: create a short link, redirect quickly, optional custom alias and expiry, click analytics. Read-heavy (about 10:1 or more), p99 redirect under about 50 ms, highly available.
- API:
POST /v1/urlsreturns the short URL;GET /{key}returns a 301 (permanent, cached by browsers, fewer hits but weaker analytics) or 302 (temporary, every click reaches you). - Key generation options: base62-encode a unique 64-bit ID from a distributed ID generator (Snowflake-style) or from pre-allocated counter ranges per server (no collisions, but sequential keys are guessable, so optionally shuffle or encrypt the ID); or hash the long URL (MD5 or SHA-256) and take 7 characters, checking for collisions; or a key-generation service that pre-creates random unused keys.
- Storage: key-value store mapping key to (long URL, owner, created, expiry), sharded by key. About 91 TB over five years before replication in the worked estimate.
- Read path: CDN or edge cache, then Redis (cache-aside), then the KV store. Hot links are a tiny fraction of keys, so the hit ratio is high.
- Analytics: publish click events to a stream asynchronously (never on the redirect's critical path) and aggregate into a warehouse.
- Extras: rate limit creation, scan for malicious URLs, expire and garbage-collect old links.
POST /urls -> API -> ID generator -> base62(id) -> KV store (sharded by key)
GET /aZ3kP9q -> CDN -> LB -> redirect service -> Redis --miss--> KV store
|
+--> click event -> Kafka -> analytics
Rate limiter (distributed, at the API gateway)
- Requirements: limit requests per user, API key or IP per rule; low added latency (about a millisecond); accurate enough; highly available.
- Algorithm: token bucket for burst tolerance, or sliding window counter for accuracy with O(1) memory.
- State: Redis keyed by
limiter:{client}:{rule}, updated atomically with a Lua script (read tokens, refill by elapsed time, decrement, set expiry) so concurrent gateway nodes do not race. - Rules stored in configuration and cached locally on each gateway node.
- Response: HTTP 429 with
Retry-AfterandX-RateLimit-Remainingheaders. - Failure policy: if Redis is unreachable, fail open (allow traffic, protect availability) for most APIs, or fail closed for expensive or abuse-prone ones; say which and why.
- Scale: shard Redis by client key; for extreme QPS, local in-memory buckets with periodic sync to a global count, accepting small overshoot. Multi-region: per-region limits or asynchronous global aggregation.
Distributed key-value store (Dynamo-style)
- Partitioning: consistent hashing with virtual nodes.
- Replication: each key stored on N = 3 successive nodes on the ring, across failure zones.
- Consistency knob: quorum reads and writes; R + W > N for read-after-write; lower values for speed.
- Conflict handling: version vectors to detect concurrent writes, resolved by the client or last-writer-wins.
- Failure handling: hinted handoff (a neighbour temporarily holds writes for a down node), read repair, and anti-entropy with Merkle trees to find and fix divergent ranges.
- Membership: gossip protocol with failure detection.
- Storage engine: LSM tree (write-ahead log, memtable, SSTables, compaction, Bloom filters), tombstones for deletes.
Distributed cache (Redis Cluster-like)
- Keys hashed into slots (Redis Cluster uses 16,384 hash slots) assigned to primary nodes, each with replicas for failover.
- Clients cache the slot map and are redirected when slots move during resharding.
- Eviction per node (LRU or LFU approximations), TTLs, memory limits.
- Hot keys: client-side local caching, key replication with suffixes, or read from replicas.
- Consistency is best effort: asynchronous replication can lose recent writes on failover, which is acceptable for a cache, not for a source of truth.
Notification system (push, SMS, email)
- Flow: producer services call a notification API (or publish events) with an idempotency key; the service validates, checks user preferences, quiet hours and opt-outs, renders templates, and enqueues per channel.
- Channel workers call provider adapters: APNs (iOS), FCM (Android), SMS gateways, email providers. Separate queues per channel and priority so a slow email provider does not delay urgent push messages.
- Device registry sharded by user or device ID, storing tokens; remove tokens the provider reports as invalid.
- Reliability: at-least-once with dedupe on the idempotency key, retries with backoff, dead-letter queue, provider failover.
- Protection: per-user and per-channel rate limits (avoid spamming), global throttles to respect provider quotas.
- Tracking: delivery, open and click events flow back into analytics; metrics on delivery rate per channel.
Services --> Notification API --> prefs/templates --> per-channel queues
| | |
push SMS email workers
(APNs/FCM) (gateway) (provider)
\_______|________/
retries, DLQ, delivery events
Distributed job scheduler / cron
- Jobs stored in a database with the next run time; an index on
next_run_at. - Scheduler nodes poll for due jobs in time buckets, claiming each atomically (conditional update or lease) so only one node dispatches it; leader election is an alternative.
- Dispatched jobs go on a queue to worker pools; workers heartbeat, and a job whose lease expires is retried.
- Jobs must be idempotent because at-least-once execution is the realistic guarantee.
- Support retries with backoff, timeouts, priorities, dependencies (a DAG, like a workflow engine), and history for auditing.
- Scale by sharding jobs across scheduler partitions; use a timing wheel or delay queue for high volumes of short delays.
Blueprints II: product designs
Product-style designs add large fan-out, real-time connections, media pipelines and geospatial data. The same levers apply, but each has one signature problem to solve.
A postal service. Chat is registered mail with delivery receipts: every letter is tracked until the recipient signs for it, and held at the post office if they are away (offline queue). A news feed is a newspaper: you can print a personalised copy for every subscriber in advance (fan-out on write) or assemble it when they walk into the shop (fan-out on read), and for a celebrity columnist with millions of readers you do the second. File sync is shipping a house move in numbered boxes so only the boxes that changed need to be re-sent (chunking and deduplication).
Chat (WhatsApp, Messenger)
- Requirements: one-to-one and group messages, delivery and read receipts, online presence, offline delivery, message history, media, optional end-to-end encryption. Low latency (under a few hundred ms), messages never lost, ordered per conversation.
- Connections: clients keep a long-lived WebSocket (or MQTT) connection to a gateway fleet; a session registry maps user ID to gateway node.
- Send path: the gateway forwards to the chat service, which assigns a per-conversation sequence number, persists first, acknowledges the sender, then routes to the recipient's gateway, or to a push notification if offline.
- Storage: messages in a wide-column store partitioned by
conversation_idand clustered by sequence number (fast "last 50 messages"). Clients sync by "give me everything after sequence N". - Groups: fan out per member via a queue; for very large groups or channels, recipients pull instead.
- Receipts: sent (server persisted), delivered (recipient device acknowledged), read (recipient opened). Client-generated message IDs make resends idempotent.
- Presence: heartbeats with a TTL in a fast store; publish presence changes only to interested contacts, and throttle.
- End-to-end encryption: the Signal protocol; servers route ciphertext only, which limits server-side search and moderation.
Sender --WebSocket--> Gateway A --> Chat service --persist--> Message store (by conversation_id)
| ack sender
+--lookup--> Session registry (user -> gateway)
+--online--> Gateway B --WebSocket--> Recipient
+--offline-> Push service (APNs / FCM)
News feed (Twitter/X, Instagram)
- Fan-out on write (push): when a user posts, write the post ID into each follower's precomputed timeline in a cache. Very fast reads; expensive writes for accounts with millions of followers.
- Fan-out on read (pull): build the timeline at read time by merging recent posts of everyone you follow. Cheap writes; slow reads.
- Hybrid (the standard answer): push for ordinary users, pull for celebrities, and merge at read time. Skip fan-out to inactive users.
- Storage: posts in a sharded store keyed by post ID; social graph in its own store; timelines as capped lists of post IDs in Redis; media on a CDN.
- Ranking: candidate generation, then a scoring model (recency, affinity, engagement), then filtering; cursor-based pagination.
- Asynchrony: fan-out workers consume "post created" events from a queue; feeds are eventually consistent, which is acceptable.
File storage and sync (Dropbox, Google Drive)
- Split files into chunks (about 4 MB); identify each by its content hash. Upload only chunks the server does not already have (deduplication and delta sync), resumable on flaky networks.
- Metadata service (relational DB, sharded by user or namespace): files, folders, versions, chunk lists, sharing permissions. Chunks live in object storage.
- Clients upload directly to object storage with pre-signed URLs, keeping large payloads off the app servers.
- Sync: each change bumps a version in a per-user change log; other devices long-poll or subscribe for notifications and pull changes since their cursor.
- Conflicts: if two devices edit the same version, keep both ("conflicted copy") or merge for structured documents.
- Durability through object-store replication or erasure coding; cold data moves to cheaper storage tiers.
Video streaming (YouTube, Netflix)
- Upload: resumable, chunked upload to object storage; a "video uploaded" event starts the pipeline.
- Transcoding pipeline: a DAG of jobs on a queue: split into segments, transcode each into multiple resolutions and bitrates in parallel, generate thumbnails, run content checks, then package as HLS or DASH manifests.
- Delivery: CDNs serve segments; adaptive bitrate players pick the quality segment by segment based on measured bandwidth. Popular content is pre-positioned at edges.
- Metadata (titles, views, comments) in databases and caches; view counts aggregated asynchronously.
- Cost dominates: storage of many renditions and CDN egress, so encode popular titles more aggressively and archive the long tail.
Ride-hailing (Uber) and proximity search
- Drivers send location every few seconds, a very high write rate. Keep current locations in memory in a geospatial index: geohash cells, a quadtree, or S2 or H3 cells, sharded by region.
- "Nearby drivers" queries the rider's cell and its neighbours, ranks candidates by estimated time of arrival, and offers the trip to one driver at a time with a timeout.
- The trip is an explicit state machine (requested, matched, arriving, in progress, completed, cancelled) stored durably; pricing includes surge based on supply and demand per cell.
- Location history streams to a log for analytics, ETA models and fraud detection. Eventual consistency is fine for positions; strong consistency is needed for assignment (a driver must not get two trips), using a conditional update or a lock.
- Proximity services such as Yelp or "find restaurants near me" use the same geospatial indexes, but the data is mostly static, so precomputed indexes and caching dominate.
Typeahead / autocomplete
- A trie of popular queries with the top-k completions cached at each prefix node, built offline from query logs (for example hourly) and served from memory.
- Shard by prefix range; replicate for read throughput; cache hot prefixes at the edge.
- Clients debounce keystrokes (about 100-200 ms) and cache recent results locally.
- Filter offensive or unsafe suggestions; personalise by blending global and user history.
Web crawler
- A URL frontier (prioritised queues) feeds fetcher workers; per-host queues enforce politeness (rate limits and
robots.txt). - Parse pages, extract links, normalise URLs, and dedupe with a Bloom filter or URL-hash store; detect near-duplicate content with content fingerprints (SimHash).
- Store raw pages in object storage; feed an indexing pipeline that builds an inverted index.
- Handle traps (infinite calendars), DNS caching, and recrawl scheduling by page change rate.
More prompts to practise
| Prompt | Signature problem |
|---|---|
| Payment system | Exactly-once effect with idempotency keys, double-entry ledger, reconciliation with the payment provider, strong consistency |
| Ticket booking / flash sale | Contention on limited inventory: holds with expiry, conditional updates, queueing users, preventing overselling |
| Leaderboard | Sorted sets (skip lists) for rank queries; sharding by score range or approximate global ranks |
| Metrics and logging platform | High-volume ingestion, time-series storage, downsampling, retention tiers |
| Search engine | Inverted index, sharding by document, scatter-gather queries, ranking |
| Collaborative editor | Operational transforms or CRDTs for concurrent edits, presence, real-time sync |
| Ad click aggregation | Stream processing with windowing, exactly-once counting, late events, reconciliation |
Mobile and embedded system design
A phone, watch or car head unit is itself a distributed system: apps, framework services, hardware abstraction layers (HALs), a modem and other co-processors talk over IPC, can crash independently, and share scarce resources. Fleets of such devices talk to cloud backends over unreliable networks. Device-side designs are judged on the usual criteria plus power, memory, intermittent connectivity, privacy and the fact that you cannot easily patch a broken device in the field.
Running a fleet of ships instead of a building. A data-centre server is a building with mains power, a fast network and an engineer down the hall. A device is a ship at sea: limited fuel (battery), a small hold (memory and storage), patchy radio contact (connectivity), and a crew that must cope alone when something breaks. You send updates only when the ship is in port and on shore power (Wi-Fi and charging), you keep a spare engine so a bad repair can be undone (A/B partitions and rollback), and you collect the logbooks in batches rather than radioing every entry (batched telemetry).
On-device constraints that change the design
| Constraint | Design response |
|---|---|
| Battery and power | Batch work, coalesce wakeups, use push (FCM) instead of polling, defer uploads to Wi-Fi and charging, let hardware sleep (radio idle states, DRX), hold wakelocks only with timeouts |
| Thermal limits | Throttle background work, schedule heavy jobs (on-device ML, indexing) when idle and cool |
| Memory and CPU | Bounded buffers and caches, streaming instead of loading whole files, low-memory-killer awareness, avoid work on the UI thread |
| Intermittent, metered networks | Offline-first local store with a sync queue, retries with backoff and jitter, resumable transfers, compression, respect metered data |
| Storage wear and space | Ring buffers and log rotation, quotas, avoid write amplification on flash |
| Privacy and security | Minimise and aggregate data on device, scrub personal data before upload, user consent, signed and verified code and configuration, hardware-backed keys |
| Heterogeneous fleet | Many OS versions, chipsets, carriers and regions; version every API and config; target rollouts by device fingerprint |
| Cannot patch instantly | Feature flags, server-driven configuration, staged rollouts, kill switches, automatic rollback |
Distributed-systems patterns inside the phone's telephony stack
The Android telephony stack (apps, framework, RIL, radio HAL, vendor modem software) is a real-time, multi-actor, partially unreliable system inside one device. Its patterns map one-to-one onto general system-design vocabulary:
| Pattern | Where it appears on the device | General system-design equivalent |
|---|---|---|
| Layering with a stable HAL | App, framework, RIL, HAL, vendor modem | Versioned service interfaces that let layers evolve independently |
| Async request/response with correlation IDs | RIL request serial numbers matched to HAL responses | Request IDs in any asynchronous RPC or messaging system |
| Publish/subscribe | Registrant lists, telephony registry callbacks | Event bus fan-out decoupling producers and consumers |
| Single-threaded event loop | Handler and Looper per component | Actor model: serialise state changes without locks |
| Explicit state machines | Service-state, call and data-connection trackers | Deterministic handling of complex transitions (trip or order state machines) |
| Death detection and recovery | HAL death recipients, flushing in-flight requests, watchdog restart | Health checks, supervisors and circuit breakers |
| Backoff and throttling | Network-mandated retry timers, data retry manager | Exponential backoff with jitter, congestion control |
| Caching with invalidation | Carrier configuration, cached service state, SIM records | Cache with invalidation on a change event |
| Rate limiting and coalescing | Throttled signal-strength and unsolicited indications | Debouncing and protecting consumers from event storms |
| Resource arbitration | Which SIM carries data on a dual-SIM phone, shared modem radio | Fair and safe scheduling of a scarce shared resource |
| Power-aware design | Wakelocks with timeouts, hardware offload, radio sleep cycles | Batching and treating energy as a first-class cost |
Blueprint: OTA (over-the-air) update system for 100 million devices
- Build and sign: the build system produces full and delta (binary diff) payloads per source-to-target build fingerprint; payloads and manifests are signed, and devices verify signatures and hashes before applying.
- Device-side safety: A/B (seamless) partitions: install to the inactive slot in the background while the user keeps working, switch slots on reboot, and automatically fall back to the old slot if the new one fails to boot (a boot-success marker). Anti-rollback counters stop downgrade attacks. Virtual A/B with snapshots saves storage.
- Rollout control: staged percentages (for example 1, 10, 50, 100 percent) targeted by model, region, carrier and build; automatic halt if crash rate, boot failures or install errors exceed the baseline; a kill switch via a signed manifest update.
- Distribution: devices check in periodically (with jitter to avoid thundering herds) with their fingerprint; the update server answers from a policy engine; payloads come from a CDN with resumable, range-based downloads, preferring Wi-Fi and charging, respecting user-chosen install windows.
- Telemetry: each device reports download, verify, install and first-boot outcomes, so the fleet dashboard shows success rate per build, model and carrier.
- Recovery: recovery mode and a factory-reset path as last resorts; staged rollouts limit the blast radius so they are rarely needed.
Build system --sign--> Payload store --> CDN
|
Update policy server <-- check-in (fingerprint, region, carrier)
| staged %, targeting, kill switch
v
Device: download (resumable, Wi-Fi) -> verify signature -> write inactive slot (B)
-> reboot into B -> boot OK? mark successful : fall back to A
-> report outcome --> telemetry --> rollout health gate (auto-halt)
Blueprint: fleet telemetry and metrics pipeline
- On device: record events (call drops, registration failures, crashes, battery drain) into a bounded ring buffer; aggregate and sample locally rather than shipping every event; scrub personal data and apply consent.
- Upload: batch uploads opportunistically (Wi-Fi and charging, or piggybacked on other network activity), compressed, with backoff on failure and a cap on data used.
- Ingest: regional collectors write to a stream; schema-versioned events; drop or quarantine malformed data.
- Processing: stream aggregation for near-real-time dashboards; batch jobs into a warehouse partitioned by build, region, carrier and model.
- Analysis: detect regressions by comparing a new build against a baseline cohort with statistical significance, so a handful of noisy devices is not mistaken for a fleet-wide bug; trigger alerts or halt rollouts automatically.
More device-flavoured prompts
Carrier configuration delivery
Layered config: platform defaults, then carrier overrides by network code, then server-pushed overrides, then a runtime cache. Versioned and signed; invalidated on SIM swap, roaming or OTA; staged rollout with a kill switch, like any feature-flag system.
Multi-SIM data-switch arbiter
Inputs: user's default data SIM, active calls, app needs, roaming and cost policy, signal quality. A state machine with hysteresis and debouncing to avoid flapping; an idempotent, transactional switch that rolls back if the new SIM fails to attach; metrics on switch latency and flaps.
Resilient framework-to-HAL channel
Death notifications fail all in-flight requests promptly instead of hanging callers; timeouts on every request (and on any wakelock held for it); idempotent re-initialisation with bounded backoff; cap outstanding requests and coalesce duplicate polls for backpressure.
Diagnostic log collection
Circular on-device buffers with triggers (capture the minutes around a failure), user consent, PII scrubbing, compression, upload on Wi-Fi, server-side indexing by device, build and failure signature for triage.
Presence and typing indicators
Ephemeral state with TTLs, heavy coalescing and throttling, fan-out only to active conversations; losing an update is acceptable, draining the battery is not.
Defect triage dashboard
Ingest issues and code-review events from trackers, deduplicate by crash signature, aggregate in near real time, and serve role-based views; mostly a data-pipeline and query-latency problem.
Mobile client architecture
- Offline-first: the local database is the UI's source of truth; a sync engine reconciles with the server using change cursors and conflict rules.
- Network layer: one shared client with connection reuse (HTTP/2 or HTTP/3), request deduplication, retries with backoff, and caching headers.
- Background work: OS job schedulers (WorkManager on Android) with constraints such as unmetered network and charging.
- Push, not poll: a push message nudges the app to sync, saving battery and server load.
- Images and media: memory and disk caches, downsampling to display size, progressive loading.
- Server-driven configuration and feature flags to change behaviour without an app release.
Senior and principal expectations: trade-off talk
At senior, staff and principal levels, the question is less "can you produce a working design?" and more "can you lead a design discussion the way you would lead a real architecture review?" Interviewers look for judgement: knowing which problems matter, quantifying them, choosing deliberately, and planning for failure, operations and evolution.
The difference between a good cook and a head chef. A good cook follows recipes well. A head chef designs the menu around the kitchen's capacity, the budget, supplier reliability and what happens when a fryer breaks mid-service, and can explain to the owner why they chose one option over another. In the interview, the recipe is the blueprint, and the head chef's reasoning is your trade-off talk: requirements, constraints, alternatives, costs and failure plans.
What changes by level
| Level | Expected |
|---|---|
| Mid-level | A working high-level design with the right building blocks; answers the interviewer's probes correctly |
| Senior | Drives the interview; clarifies requirements and quantifies; goes deep on one or two areas; handles failure modes; justifies technology choices |
| Staff | Identifies the few decisions that matter most; compares real alternatives with costs; covers operability, migration from the current system, multi-region, security and cost; thinks about team boundaries |
| Principal | Frames the problem in business terms; sets principles and non-negotiables; plans an evolution roadmap; balances org, platform and long-term maintainability; spots cross-system risks |
How to talk about trade-offs
- Name the alternatives: "We could fan out on write or on read. Write gives fast reads but costs millions of writes per celebrity post; read is cheap to write but slow to render."
- Tie the choice to a requirement: "Since feed reads outnumber posts 100 to 1 and a 2-second delay is acceptable, I will push for normal users and pull for accounts over about 100,000 followers."
- State what you give up: "The cost is eventual consistency: a follower may see a post a few seconds late."
- Say when you would revisit: "If we need strict ordering across regions later, we would move this to a consensus-backed log."
Trade-offs you should be able to discuss fluently
| Trade-off | One side | Other side |
|---|---|---|
| Consistency vs availability and latency | Correctness for money, inventory, identity | Uptime and speed for feeds, counts, presence |
| Latency vs throughput | Interactive requests, small batches | Batching, compression, pipelines |
| SQL vs NoSQL | Transactions, joins, constraints | Horizontal write scale, flexible schema |
| Push vs pull | Low read latency, write amplification | Cheap writes, read-time work |
| Sync vs async | Simple, immediate result | Decoupled, resilient, eventually consistent |
| Normalise vs denormalise | One source of truth, cheaper writes | Fast reads, duplicated data to keep in sync |
| Build vs buy (managed service) | Control, custom needs | Speed, less operational burden, vendor lock-in |
| Monolith vs microservices | Simplicity, fast iteration for small teams | Independent scaling and deployment for large orgs |
| Accuracy vs cost | Exact counts, sliding logs | Approximate sketches, sampling, HyperLogLog |
Senior checklist beyond the happy path
- Failure modes: node, zone and region loss; dependency slowness; overload; bad deploys; data corruption. What is the blast radius, and how is it contained?
- Data lifecycle: retention, deletion (including privacy requests), backups and restore drills, schema migration without downtime (expand, migrate, contract).
- Security and privacy: authentication and authorisation, encryption in transit and at rest, secrets management, abuse and fraud, data residency.
- Operability: SLOs and error budgets, dashboards, alerting, on-call runbooks, capacity planning.
- Rollout and migration: feature flags, canaries, dual writes and backfills, shadow traffic, rollback plans.
- Cost: storage tiers, CDN egress, compute efficiency; call out the dominant cost driver.
- Evolution: what you would build for version 1 versus at 10x and 100x scale.
Consensus, clocks and probabilistic sketches
A few primitives show up in almost every senior design once you have more than one writer or need an approximate answer over huge streams. Know the story of each well enough to draw it, not just name it.
A committee that must agree on the next line in a shared minute book (Raft): they elect a chair, only the chair proposes the next line, and a line is official once a majority have written it down. Wall clocks in different rooms drift, so you do not order events by "what time my watch said" (logical clocks / version vectors). A Bloom filter is a receptionist who can say "we have definitely never seen this guest" after glancing at a tiny board of ticks, but who sometimes says "maybe" for a stranger; they never miss a real guest.
Raft in one picture
- Leader election: a follower that times out of heartbeats becomes a candidate, increments its term, and asks for votes. A node votes for at most one candidate per term. Majority wins; the new leader sends heartbeats.
- Log replication: clients write to the leader. The leader appends to its log and ships the entry to followers. The entry is committed once a majority has stored it; then the leader applies it and acknowledges the client.
- Safety: a candidate can win only if its log is at least as up to date as the voters' logs, so a committed entry cannot be lost. With 2f+1 nodes you tolerate f failures.
Use Raft (etcd, ZooKeeper-style, Consul, TiKV regions) for leader election, configuration, and any small amount of data that must be linearizable. Do not put the request path of a high-QPS data plane through a single Raft group.
Ordering without a global clock
- Lamport clocks: a single counter that increments on every event and is piggy-backed on messages; they give a partial order (causality) but not concurrent-event identity.
- Version vectors: one counter per replica; comparing two vectors tells you happens-before or concurrent (conflict). Used in leaderless stores to detect concurrent writes.
- TrueTime (Spanner): GPS/atomic clocks with a bounded uncertainty interval; commit waits out the uncertainty so commits have a real-time order. Rare outside that design.
Bloom filters and sketches
- Bloom filter: bit array plus k hash functions.
addsets k bits;maybe_containsis true if all k bits are set. False positives, never false negatives. ~10 bits/item for a 1 percent false-positive rate. Cannot delete (unless you use a counting Bloom filter). Classic uses: cache penetration ("this key is definitely not in the DB"), crawler URL-seen sets, LSM-tree SSTable indexes. - Count-Min Sketch: approximate frequencies for heavy hitters in a stream, with bounded overestimate, tiny memory.
- HyperLogLog: approximate distinct count (cardinality) in about 1โ2 KB.
Quick revision
- Start every design with functional and non-functional requirements, then estimate scale; never start with boxes.
- 1 million requests per day is about 12 per second; peak is about 2-3x average; a day is about 10^5 seconds.
- Memory is about 100 ns; 1 MB sequential from RAM is about 10-50 μs (not a few microseconds). SSD random is about 100 μs, same-DC round trip about 0.5 ms, cross-continent about 150 ms.
- Scale out with stateless services behind load balancers; keep session state in a shared store or token.
- Report latency as percentiles; fan-out amplifies tail latency.
- 99.9% availability is about 8.8 hours of downtime per year; 99.99% is about 53 minutes.
- Serial dependencies multiply availability down; independent redundancy multiplies failure probability down.
- CAP: during a partition choose consistency or availability; PACELC adds latency vs consistency in normal operation.
- Quorum reads see the latest write when R + W > N.
- Pick consistency per feature: strong for money and inventory, eventual for feeds and counts.
- L4 load balancers route by IP and port; L7 by HTTP content, with TLS termination and path routing.
- CDNs cache static and cacheable content at the edge; use content-hashed file names instead of purging.
- Cache-aside is the default; on writes, update the database then delete the cache key.
- Prevent cache stampedes with request coalescing, TTL jitter and stale-while-revalidate.
- LRU is O(1) with a hash map plus a doubly linked list.
- Choose the database from access patterns: relational for transactions and joins, wide-column for massive writes, KV for key lookups.
- Indexes speed reads to O(log n) but slow writes and use storage; composite indexes follow the leftmost prefix.
- Replication scales reads and availability; async replication means lag and possible data loss on failover.
- Sharding scales writes and storage; the shard key must spread load and match the main query.
- Consistent hashing moves only about 1/N of keys when a node is added; virtual nodes balance load.
- Queues decouple, buffer spikes and enable retries; they give at-least-once delivery and per-partition ordering.
- Exactly-once is at-least-once delivery plus idempotent processing.
- Idempotency keys make retried writes safe; the outbox pattern solves the dual-write problem.
- Token bucket allows bursts with an average rate; sliding window counter is accurate with O(1) memory.
- Rate limiter state lives in Redis updated atomically; return 429 with Retry-After.
- Every remote call needs a timeout; retry only idempotent calls, with exponential backoff and jitter.
- Circuit breakers fail fast on a sick dependency; bulkheads isolate resource pools.
- URL shortener: base62 of a unique ID, read-heavy KV store, cache and CDN, async analytics.
- News feed: hybrid fan-out, push for normal users and pull for celebrities.
- Chat: WebSocket gateways, persist before acknowledging, per-conversation sequence numbers, push when offline.
- File sync: content-hashed chunks, deduplication, metadata DB, pre-signed uploads, change cursors.
- Ride-hailing: in-memory geospatial index (geohash, quadtree, S2/H3) and exclusive trip assignment.
- OTA: signed payloads, A/B partitions with automatic fallback, staged rollout gated on telemetry, kill switch.
- Device design treats battery, memory, connectivity and privacy as first-class constraints; batch and defer work.
- Raft: elect a leader, replicate a log, commit on majority; 2f+1 nodes tolerate f failures. Keep it off the high-QPS data path.
- Bloom filter: maybe-present, never false-negative, no deletes; ~10 bits/item at 1% false positives. HyperLogLog for distinct counts, Count-Min for heavy hitters.
- Version vectors detect concurrent writes; do not order distributed events by wall-clock timestamps alone.
- Senior answers name alternatives, tie choices to requirements, state what is given up, and cover failure, operations and cost.
Glossary
- A/B partitions
- Two copies of system partitions so an update installs to the inactive slot and the device can fall back if it fails to boot.
- ACID
- Atomicity, consistency, isolation, durability: the transaction guarantees of relational databases.
- Adaptive bitrate
- Streaming technique where the player switches video quality segment by segment based on bandwidth (HLS, DASH).
- API gateway
- Entry point that handles authentication, rate limiting, routing and other cross-cutting concerns for backend services.
- Availability
- Fraction of time a system serves requests successfully, often expressed in nines.
- Backpressure
- Signalling producers to slow down when consumers cannot keep up.
- BASE
- Basically available, soft state, eventual consistency: the looser guarantees of many NoSQL systems.
- Bloom filter
- Compact probabilistic set that can say "definitely not present" or "probably present".
- Bulkhead
- Isolating resources (thread or connection pools) per dependency so one failure cannot exhaust all of them.
- Cache-aside
- Caching strategy where the application reads the cache and, on a miss, loads from the database and fills the cache.
- CAP theorem
- During a network partition a distributed system must choose between consistency and availability.
- CDN
- Content delivery network: edge servers that cache content close to users.
- Change data capture
- Streaming a database's committed changes (from its log) to other systems.
- Circuit breaker
- Pattern that stops calling a failing dependency for a while and fails fast, then probes for recovery.
- Consensus
- Protocol (Raft, Paxos) for a group of nodes to agree on a value or leader despite failures.
- Consistent hashing
- Mapping keys and nodes onto a ring so adding or removing a node moves only a small fraction of keys.
- CQRS
- Command query responsibility segregation: separate models for writes and for reads.
- CRDT
- Conflict-free replicated data type that merges concurrent updates automatically and deterministically.
- Dead-letter queue
- Queue that holds messages that repeatedly failed processing, for inspection.
- Denormalisation
- Storing redundant copies of data to make reads faster at the cost of more complex writes.
- Error budget
- The allowed amount of unreliability implied by an SLO, spent on launches and risk.
- Eventual consistency
- Guarantee that replicas converge to the same value if writes stop.
- Fan-out
- Delivering one event or request to many recipients or backends.
- Fencing token
- Monotonically increasing number issued with a lock so storage can reject writes from stale lock holders.
- Geohash
- Encoding of latitude and longitude into a string where shared prefixes mean nearby locations.
- Hinted handoff
- A replica temporarily accepts writes for an unavailable node and forwards them when it returns.
- Hot key
- A single key receiving a disproportionate share of traffic, overloading one shard or cache node.
- Idempotency
- Property that repeating an operation has the same effect as performing it once.
- Leader-follower replication
- Writes go to one leader that streams changes to read-only followers.
- Linearizability
- Strongest single-object consistency: operations appear instantaneous and in real-time order.
- Load balancer
- Component that distributes requests across healthy servers.
- LSM tree
- Log-structured merge tree: write-optimised storage that flushes sorted files and compacts them in the background.
- Outbox pattern
- Writing events to a table in the same transaction as business data, then publishing them reliably.
- PACELC
- Extension of CAP: under partition choose availability or consistency, else choose latency or consistency.
- Quorum
- Minimum number of replicas that must respond to a read or write; R + W > N gives overlapping reads and writes.
- Rate limiting
- Restricting how many requests a client may make in a time period.
- Replication lag
- Delay between a write on the leader and its appearance on followers.
- Saga
- A long-running transaction split into local steps with compensating actions on failure.
- Sharding
- Partitioning data across nodes by a key so each node holds a subset.
- SLA, SLO, SLI
- Contractual promise, internal target, and measured indicator of service quality.
- Snowflake ID
- 64-bit roughly time-ordered unique ID built from a timestamp, machine ID and sequence number.
- Split brain
- Failure where two nodes both believe they are the leader and accept conflicting writes.
- Stampede
- Many requests missing the cache at once and overloading the backing store (thundering herd).
- Tail latency
- The slowest requests, measured by high percentiles such as p99.
- Token bucket
- Rate-limiting algorithm where tokens refill at a steady rate and each request consumes one, allowing bursts.
- Raft
- Consensus protocol: a leader replicates a log; an entry is committed when a majority has stored it.
- Version vector
- Per-replica counters used to tell whether two updates happen-before each other or conflict.
- HyperLogLog
- Tiny probabilistic structure that estimates the number of distinct items in a stream.
- Write-ahead log
- Append-only durable log of changes written before in-memory state is treated as committed; replayed after a crash.
Interview questions
Fundamentals
What is the difference between vertical and horizontal scaling?
Vertical scaling (scale up) means a bigger machine: more CPU, memory or faster disks. It is simple and needs no code changes, but has a hard ceiling, rising cost, and remains a single point of failure. Horizontal scaling (scale out) means more machines behind a load balancer. It grows almost without limit and tolerates failures, but requires stateless services, data partitioning and handling of network and coordination complexity. Large systems scale out, though scaling a database up first is often the pragmatic early step.
Why should application servers be stateless?
If a server keeps no per-user state between requests, any server can handle any request. That lets you add or remove servers freely (autoscaling), replace a failed server without losing sessions, deploy with rolling restarts, and balance load evenly without sticky sessions. State moves to shared stores built for it: sessions in Redis or a database, or a signed token such as a JWT carried by the client; files in object storage.
What is the difference between latency and throughput?
Latency is the time one request takes; throughput is how many requests (or bytes) the system handles per unit time. They are related but distinct: batching raises throughput but can raise latency, and a system can have high throughput with poor latency. By Little's law, concurrency = throughput x latency, so for a fixed concurrency limit, lower latency means higher throughput. Optimise the one the product needs: interactive APIs care about latency, batch pipelines about throughput.
Why do we measure latency with percentiles instead of the average?
Latency distributions have long tails; the average hides the slow requests that users actually notice. p50 is the typical experience, p99 is the experience of 1 in 100 requests, which for a heavy user can mean several slow requests per session. When one page fans out to many backends, the page is as slow as its slowest call, so backend p99 becomes user-facing median. SLOs are usually written on p95 or p99.
What does 99.99% availability mean in practice?
About 53 minutes of downtime per year, or about 4.4 minutes per month. For comparison, 99.9% is about 8.8 hours per year and 99.999% about 5 minutes. Serial dependencies multiply: a service at 99.99% that depends on a database at 99.99% is about 99.98%. Achieving four nines requires redundancy with no single points of failure, automated failover, safe deployment practices, and often multiple zones.
What is a single point of failure and how do you remove it?
A single point of failure (SPOF) is any component whose failure brings down the whole system: one database, one load balancer, one region, one DNS provider, even one engineer's credentials. Remove it with redundancy (multiple instances across zones), health checks with automatic failover (leader election for databases), replication of data, and avoiding shared mutable singletons. Test failover regularly, since untested redundancy often fails when needed.
Explain the CAP theorem with a real example.
When the network partitions, a distributed system must choose between consistency (every read returns the latest write) and availability (every request to a live node gets a successful response). A bank ledger chooses CP: the minority side rejects writes rather than let balances diverge. A shopping cart or social feed chooses AP: both sides keep accepting writes and reconcile later. Without a partition you can have both; CAP only forces the choice during partitions, which cannot be ruled out in real networks.
Strong vs eventual consistency: how do you choose?
Strong consistency means every read reflects the latest committed write; it is needed for balances, inventory, uniqueness constraints (usernames) and authorisation changes. Eventual consistency means replicas converge over time; it is fine for likes, view counts, feeds, recommendations and presence, and gives better latency and availability at scale. Choose per data type within the same system, and use in-between guarantees such as read-your-writes where users would otherwise notice.
SQL vs NoSQL: how do you decide?
Start from access patterns and consistency needs. Choose relational (PostgreSQL, MySQL) for structured data with relationships, ad-hoc queries, joins, constraints and ACID transactions (orders, payments, accounts); it scales further than people assume with indexes, replicas and partitioning. Choose NoSQL when you need massive horizontal write throughput, simple key-based access, flexible schemas, or a specific model: key-value (sessions), wide-column (messages, time series), document (catalogues), graph (social connections). Many systems use both.
What is a database index and what does it cost?
An index is an auxiliary data structure, usually a B+ tree, that keeps a column's values sorted with pointers to rows, turning a full table scan (O(n)) into an O(log n) lookup and enabling efficient range queries and sorting. Costs: every insert, update and delete must maintain each index (slower writes), extra storage, and the optimiser can pick badly with too many. Index columns used in WHERE, JOIN and ORDER BY, respect the leftmost-prefix rule for composite indexes, and avoid over-indexing write-heavy tables.
What is the difference between replication and sharding?
Replication copies the same data to multiple nodes, improving read throughput, availability and durability, but not write capacity (every replica still applies every write). Sharding splits different data across nodes by a key, scaling writes and storage, but adds cross-shard query complexity. Real systems combine both: each shard is replicated.
What does a load balancer do, and which algorithms can it use?
It distributes incoming requests across a pool of servers, runs health checks to remove unhealthy ones, and often terminates TLS. Algorithms: round robin, weighted round robin for mixed capacity, least connections for long or uneven requests (WebSockets), least response time, IP or consistent hashing for affinity and cache locality, and power-of-two-choices (pick two random servers, choose the less loaded). The balancer itself must be redundant (active-passive pair or a managed service).
L4 vs L7 load balancing?
Layer 4 balancers route by IP address and port without reading the payload: very fast, protocol agnostic, suitable for any TCP or UDP traffic. Layer 7 balancers understand HTTP: they route by path, host, header or cookie, terminate TLS, rewrite headers, retry, compress and support sticky sessions, at a bit more CPU cost. A common setup is L4 at the edge for raw scale, then L7 for application routing.
What is a CDN and when does it help?
A content delivery network caches content on edge servers near users. It cuts latency (fewer long round trips), offloads the origin, absorbs traffic spikes and DDoS, and can terminate TLS near users. It helps most for static and cacheable content (images, video segments, JavaScript, CSS, public API responses) and globally distributed users. It helps little for personalised, rapidly changing or write-heavy traffic. Use versioned file names for cache busting and signed URLs for private content.
What is cache-aside, and why is it the most common strategy?
The application checks the cache first; on a miss it reads the database, stores the result in the cache with a TTL, and returns it. On a write it updates the database and deletes the cache entry. It is common because the cache only holds data that is actually requested, the application controls the logic, and if the cache fails the system still works (just slower). Downsides: the first request after a miss is slow, and there are brief windows of staleness.
What are common cache eviction policies?
LRU evicts the least recently used item and suits workloads where recent use predicts future use; it is O(1) with a hash map and a doubly linked list. LFU evicts the least frequently used item and suits stable popularity but needs ageing. FIFO and random are simple and sometimes adequate. TTL-based expiry bounds staleness independently of memory pressure. Redis offers approximate LRU and LFU policies that sample keys rather than tracking exact order.
Why would you introduce a message queue?
To decouple producers from consumers (neither needs the other to be up), to absorb traffic spikes and level load, to move slow work off the request path (emails, image processing), to fan events out to several independent consumers, and to retry failed work durably. The cost is eventual consistency, operational overhead, and the need for idempotent consumers and careful ordering design.
What is idempotency and why does it matter in distributed systems?
An operation is idempotent if performing it several times has the same effect as once. Networks time out, clients retry, and queues redeliver, so duplicates are normal. Without idempotency, a retried payment charges twice or a redelivered message creates two orders. Achieve it with client-generated idempotency keys stored by the server, unique constraints on natural keys, conditional or versioned updates, and deduplication tables in consumers. HTTP GET, PUT and DELETE are defined as idempotent; POST is not.
REST vs gRPC vs GraphQL?
REST (HTTP and JSON) is simple, universal, cacheable and ideal for public APIs. gRPC uses HTTP/2 and Protocol Buffers: compact, fast, strongly typed with generated clients and streaming, ideal for internal service-to-service calls, but less friendly to browsers and humans. GraphQL lets clients request exactly the fields they need in one round trip, great for mobile clients aggregating many resources, but caching, authorisation per field and query cost control are harder.
Offset vs cursor pagination?
Offset pagination (?page=5&size=20) is simple and supports jumping to a page, but the database must skip all earlier rows (slow for deep pages), and inserts or deletes between requests cause skipped or duplicated items. Cursor (keyset) pagination returns an opaque cursor encoding the last seen sort key; the next query is "items after this key", which uses an index, stays fast at any depth, and is stable under concurrent writes. Use cursors for feeds and large lists.
What is rate limiting and where do you apply it?
Rate limiting caps how many requests a client may make per time window, protecting services from abuse, runaway clients and overload, and enforcing fair use or pricing tiers. Apply it at the API gateway or edge (per user, API key or IP), within services for expensive operations, and on the client side to be a good citizen. Common algorithms are token bucket, leaky bucket, fixed window, sliding window log and sliding window counter. Rejected requests get HTTP 429 with a Retry-After header.
What are ACID and BASE?
ACID describes transactional guarantees in relational databases: atomicity (all or nothing), consistency (constraints hold), isolation (concurrent transactions behave as if serial, depending on the isolation level) and durability (committed data survives crashes, via a write-ahead log). BASE describes many large NoSQL systems: basically available, soft state, eventually consistent. ACID favours correctness; BASE favours availability and scale.
Monolith vs microservices: when would you choose each?
A monolith is one deployable unit: simple to develop, test, debug and deploy, with in-process calls and easy transactions; ideal for small teams and new products. Microservices split the system into independently deployable services owned by separate teams: independent scaling and releases, fault isolation and technology freedom, at the cost of network calls, partial failures, distributed data and a large observability and platform investment. A common path is a well-modularised monolith first, then extracting services where team or scaling boundaries demand it.
What are SLIs, SLOs, SLAs and error budgets?
An SLI (service level indicator) is a measurement, such as the percentage of requests served successfully under 300 ms. An SLO (objective) is the internal target for that SLI, such as 99.9% over 30 days. An SLA (agreement) is an external contract with consequences, usually looser than the SLO. The error budget is 100% minus the SLO; while budget remains, teams can ship risky changes, and when it is exhausted they focus on reliability.
What is the difference between a forward proxy and a reverse proxy?
A forward proxy acts on behalf of clients: clients send requests through it to reach the internet (corporate egress filtering, anonymity, caching). A reverse proxy acts on behalf of servers: clients talk to it thinking it is the server, and it forwards to backend servers, providing load balancing, TLS termination, caching, compression and protection (Nginx, Envoy, HAProxy). Load balancers and API gateways are specialised reverse proxies.
How does DNS resolution work, and how is DNS used in system design?
The client asks a recursive resolver, which queries a root server, then the top-level domain server, then the domain's authoritative server, caching each answer for its TTL. In designs, DNS provides geographic or latency-based routing to the nearest region, weighted routing for gradual migrations, and failover by changing records. Low TTLs allow faster failover but increase query load, and some clients ignore TTLs, so DNS failover is never instantaneous.
Going deeper
What is a Bloom filter and when do you use one?
A Bloom filter is a bit array plus k hash functions. Adding an item sets k bits; a lookup that finds any bit unset means the item was never added (no false negatives). If all k bits are set, the item is probably present (false positives). About 10 bits per item gives roughly a 1 percent false-positive rate. You cannot delete (unless you use a counting variant) and you cannot list the items. Use it to skip work: "this URL was definitely not crawled", "this key is definitely not in the database" (cache penetration), or "this key is definitely not in this SSTable" (LSM reads). If a false positive is expensive, keep the error rate low or confirm with the source of truth.
What is a write-ahead log?
A write-ahead log (WAL) is an append-only file of changes written to durable storage before the corresponding in-memory structure (database page, memtable) is treated as committed. On crash, the store replays the log to recover. Sequential appends are fast; the log is the durability mechanism behind ACID's D and behind LSM memtable flushes. Related ideas: Kafka is a distributed WAL; the outbox pattern is a WAL of events next to business rows; Redis AOF is a WAL of commands.
What is PACELC and how does it extend CAP?
PACELC says: if there is a Partition, choose Availability or Consistency; Else, in normal operation, choose Latency or Consistency. It captures the everyday trade-off CAP ignores: synchronous replication to other replicas or regions makes every write slower, while asynchronous replication is fast but readers can see stale data. DynamoDB and Cassandra are typically PA/EL; Google Spanner and synchronously replicated relational databases are PC/EC.
Explain quorum reads and writes. Why R + W > N?
Each key is stored on N replicas. A write succeeds when W replicas acknowledge; a read queries R replicas and takes the newest version. If R + W > N, the set of replicas read must overlap the set written, so at least one replica in every read has the latest write. N=3, W=2, R=2 tolerates one failed replica for both reads and writes. W=1, R=1 is fast but may return stale data. Quorums alone are not full linearizability (concurrent writes and sloppy quorums complicate it), so systems add versioning and read repair.
Compare linearizable, causal, read-your-writes and eventual consistency.
Linearizable: the system behaves like a single copy; once a write completes, every later read sees it. Needed for locks, leader election, unique constraints; costs latency and availability. Causal: causally related operations are seen in order everywhere (a reply never appears before its post), but unrelated ones may differ; available under partitions. Read-your-writes: a user always sees their own updates, implemented by routing their reads to the leader or tracking replication position. Eventual: replicas converge if writes stop; cheapest, but readers can see old or out-of-order data.
Compare write-through, write-back and write-around caching.
Write-through: write to cache and database synchronously; the cache is always fresh, but writes are slower and the cache fills with data that may never be read. Write-back: write to the cache and flush to the database later; very fast and allows batching, but risks losing data if a cache node fails before flushing. Write-around: write only to the database; the cache fills on later reads, avoiding pollution by write-once data, at the cost of a miss on the first read. Most web systems use cache-aside reads with write-around plus invalidation.
What is a cache stampede and how do you prevent it?
When a popular key expires or is evicted, many concurrent requests miss at once and all hit the database to rebuild it, which can overload it. Prevention: request coalescing or a per-key lock so only one request rebuilds while others wait or get stale data; serving stale-while-revalidate; adding random jitter to TTLs so keys do not expire together; probabilistic early refresh of hot keys; and pre-warming caches after deploys or restarts.
How do you keep a cache consistent with the database?
Treat the database as the source of truth. On writes, commit to the database first, then delete (not update) the cache key, so the next read reloads fresh data; deleting avoids races where concurrent writers leave an older value cached. Bound staleness with TTLs. For stronger guarantees, drive invalidation from the database's change stream (change data capture) so every committed change invalidates the cache, or use versioned keys. Accept that perfect consistency between two systems without a transaction is not achievable; state how much staleness the product tolerates.
How does consistent hashing work and why is it useful?
Nodes and keys are hashed onto the same circular space; each key belongs to the first node clockwise from it. When a node joins, it takes only the keys between it and its predecessor; when a node leaves, only its keys move to its successor, about 1/N of the data, whereas hash mod N remaps almost every key when N changes. Virtual nodes (many ring positions per physical node) even out the load and allow weighting by capacity. Used in distributed caches, Dynamo-style stores, and load balancers needing affinity.
How do you choose a shard key, and how do you handle a hot shard?
A good shard key has high cardinality, distributes reads and writes evenly, and matches the dominant access pattern so most queries hit a single shard (user_id for user data, conversation_id for messages). Avoid monotonically increasing keys with range sharding (all new writes hit one shard). For hot shards: split the hot range, add a random suffix to hot keys and scatter-gather on read, cache hot entities aggressively, give celebrities special handling, or move a hot tenant to a dedicated shard via a directory service.
What is replication lag, and how do you give users read-your-writes consistency?
With asynchronous replication, followers apply the leader's changes after a delay, from milliseconds to seconds under load. A user who updates their profile and immediately reads from a follower may see the old value. Fixes: read the user's own data from the leader, or from the leader for a short window after they write; track the log position of the user's last write and only read from followers that have caught up; or pin a user to one replica for monotonic reads.
How does leader failover work, and what is split brain?
Followers or a coordinator detect that the leader has stopped heartbeating, elect a new leader (ideally the most up-to-date follower) through a consensus system, and redirect clients. Split brain happens when the old leader was only slow or partitioned, not dead, and keeps accepting writes, so two leaders diverge. Prevent it with majority-based election (a leader needs a quorum), leases that expire, and fencing tokens that storage checks to reject writes from a deposed leader. Asynchronous replication also means the newly promoted leader may be missing the last few writes.
Kafka vs RabbitMQ (log vs queue)?
Kafka is a distributed, partitioned, append-only log. Messages are retained for a configured time regardless of consumption; many consumer groups read independently by offset; ordering is per partition; throughput is very high; replay is easy. It suits event streaming, analytics pipelines, event sourcing and change data capture. RabbitMQ is a message broker with queues: messages are routed via exchanges, delivered to one consumer, and removed on acknowledgement; it supports per-message routing, priorities and delays. It suits task queues and complex routing at moderate scale.
At-most-once, at-least-once, exactly-once: what is realistic?
At-most-once acknowledges before processing and can lose messages. At-least-once acknowledges after processing and can duplicate on retries or consumer crashes; it is the practical default. True exactly-once delivery across independent systems is not possible in general, but you can achieve exactly-once effects by combining at-least-once delivery with idempotent processing (dedupe by message ID, conditional writes) or with transactions that commit the consumer offset and the result atomically (as Kafka transactions do within Kafka).
What is the dual-write problem and how does the outbox pattern solve it?
A service that writes to its database and then publishes an event can fail in between, leaving the database updated but no event sent (or an event sent for a rolled-back write). The transactional outbox writes the business change and an event row into an outbox table in the same local transaction. A separate relay (polling or change data capture) reads the outbox and publishes to the queue, retrying until it succeeds. Delivery becomes at-least-once, so consumers must be idempotent.
Saga vs two-phase commit for transactions across services?
Two-phase commit has a coordinator ask all participants to prepare, then commit; it gives atomicity but holds locks during the protocol and blocks if the coordinator fails, which hurts availability and couples services. A saga breaks the business transaction into local transactions, each publishing an event or called by an orchestrator; if a step fails, compensating transactions undo earlier steps (refund payment, release inventory). Sagas are available and loosely coupled but only eventually consistent, and compensations must be designed explicitly. Microservices generally prefer sagas.
How should retries be done safely?
Retry only transient failures (timeouts, 503s), never permanent ones (400s). Retry only idempotent operations or those protected by idempotency keys. Use exponential backoff (for example 100 ms, 200 ms, 400 ms) with random jitter so clients do not retry in lockstep, cap the attempt count and total time, and use a retry budget (for example retries at most 10 percent of requests). Retry at one layer only, to avoid multiplicative retry storms, and respect Retry-After headers.
How does a circuit breaker work?
It wraps calls to a dependency and tracks failures. In the closed state calls pass through. When failures or timeouts exceed a threshold, it opens and fails calls immediately (or serves a fallback) for a cooldown period, protecting both the caller's threads and the struggling dependency. After the cooldown it goes half-open and lets a few trial calls through; success closes it, failure reopens it. Pair it with timeouts, bulkheads and meaningful fallbacks such as cached data or a degraded feature.
How do you generate unique IDs in a distributed system?
Options: UUIDv4 (random 128-bit, no coordination, but not sortable and poor for B-tree locality; UUIDv7 adds a time prefix); Snowflake-style 64-bit IDs of timestamp + machine ID + per-millisecond sequence (sortable by time, compact, no coordination after machine ID assignment, but needs clock care); database ticket servers or range allocation where each server reserves a block of IDs (simple, few round trips). Choose by whether you need ordering, size limits, and guessability concerns.
B-tree vs LSM-tree storage engines?
B-trees (most relational databases) update pages in place: reads are fast and predictable with one tree lookup, and range scans are efficient, but random writes cause random I/O and write amplification. LSM trees (Cassandra, RocksDB) append writes to a log and an in-memory table, flush sorted immutable files, and merge them by compaction: writes are sequential and very fast, but reads may check several files (mitigated by Bloom filters) and compaction consumes I/O in the background. Choose LSM for write-heavy workloads, B-trees for read-heavy and transactional ones.
What are transaction isolation levels and the anomalies they prevent?
| Level | Prevents | Still allows |
|---|---|---|
| Read uncommitted | Almost nothing | Dirty reads |
| Read committed | Dirty reads | Non-repeatable reads, phantoms, lost updates |
| Repeatable read / snapshot | Non-repeatable reads (and phantoms in some engines) | Write skew |
| Serializable | All anomalies | Nothing, at a performance cost |
Lost updates can also be prevented with SELECT ... FOR UPDATE or optimistic version checks. Write skew example: two doctors both go off call because each saw the other still on call.
Polling vs long polling vs server-sent events vs WebSockets?
Short polling: the client asks every few seconds; simple but wasteful and laggy. Long polling: the server holds the request until there is data or a timeout; near real time over plain HTTP. Server-sent events: one long-lived HTTP response streaming server-to-client events; simple, auto-reconnect, one direction only. WebSockets: a full-duplex persistent connection; best for chat, games and collaborative editing, but stateful connections must be load-balanced, drained on deploy, and mapped to users in a session registry.
Estimate the storage needed for a photo-sharing app with 10 million daily uploads.
Assume an average stored photo of 2 MB after compression, plus about 3 smaller renditions totalling 0.5 MB, so 2.5 MB per upload. Daily: 10 M x 2.5 MB = 25 TB per day. Yearly: about 9 PB. With 3x replication that is about 27 PB per year, or about 1.5x with erasure coding (about 14 PB). Metadata is small by comparison (about 1 KB per photo, 3.6 TB per year). Conclusions: blob storage with erasure coding and cold tiers for old photos, a CDN for delivery, and metadata in a sharded database.
Estimate read and write QPS for a Twitter-like service with 200 million daily active users.
Assume each user posts 0.5 times a day and reads their timeline 20 times a day. Writes: 200 M x 0.5 / 86,400 is about 1,200 posts per second, peak about 3,000 to 5,000. Timeline reads: 200 M x 20 / 86,400 is about 46,000 per second, peak about 100,000 to 150,000. The system is read-heavy by about 40:1, so precomputed timelines in cache (fan-out on write) make sense, and the fan-out work is posts per second times average followers, which is where celebrities need special handling.
How do you implement a distributed lock correctly?
Use a system built on consensus (etcd, ZooKeeper, Consul) or a carefully designed Redis lock with a lease: acquire atomically with an expiry (SET key value NX PX 30000) and a unique owner value, release only if you still own it (compare-and-delete). Because a holder can pause (garbage collection, network) past its lease, the resource being protected must check a fencing token, a monotonically increasing number issued with each lock grant, and reject writes carrying an older token. Prefer designs that avoid locks, such as partitioning ownership or conditional writes.
What should you monitor in a production service?
The four golden signals: latency (percentiles, separately for successes and errors), traffic (requests per second), errors (rate by type), and saturation (CPU, memory, queue depth, connection pool usage). Add business metrics (orders per minute), dependency health, and for queues, consumer lag. Combine metrics with structured logs carrying request IDs and distributed traces. Alert on SLO burn rate and user-visible symptoms rather than every internal cause, and give each alert a runbook.
OLTP vs OLAP: why separate them?
OLTP databases serve the application: many small, indexed reads and writes with low latency and transactions, stored row by row. OLAP systems serve analytics: few, large queries scanning and aggregating billions of rows, stored column by column with heavy compression. Running analytics on the OLTP database competes for resources and slows the product. Instead, stream changes (change data capture or ETL) into a warehouse or lakehouse such as BigQuery, Redshift, Snowflake or ClickHouse.
How can a distributed rate limiter stay accurate without adding much latency?
A central Redis counter updated atomically with a Lua script is accurate but adds a network round trip per request and becomes a hot dependency. Options to trade accuracy for latency: local token buckets on each gateway node with a share of the global limit, periodically rebalanced; batching increments asynchronously to the central store and enforcing slightly conservative local limits; sharding Redis by client key; and sticky routing of a client to the same gateway so local state suffices. Decide in advance whether to fail open or closed if the store is unavailable.
Advanced
How does Raft elect a leader and commit a write?
Each node is follower, candidate or leader. Followers expect heartbeats; on timeout a follower increments its term, becomes a candidate, and requests votes. A node votes for at most one candidate per term; majority wins and the leader starts heartbeats. Clients write to the leader, which appends to its log and replicates the entry. The entry is committed when a majority of nodes have stored it; the leader then applies it and replies. Safety: a candidate cannot win unless its log is at least as fresh as the voters', so committed entries are never overwritten. With 2f+1 nodes you tolerate f failures. Use Raft for metadata, leader election and small linearizable state (etcd, ZooKeeper-like systems), not as the data plane for high-QPS records.
Lamport clocks vs version vectors vs wall clocks?
Wall clocks drift and jump; do not order events across machines by timestamp alone (last-writer-wins can silently drop an update). Lamport clocks are a single counter incremented on every event and sent with messages; they preserve causality (if A happened-before B then L(A) < L(B)) but two equal counters can still be concurrent. Version vectors keep one counter per replica; comparing two vectors tells you happens-before or concurrent (a conflict to merge). Leaderless stores use version vectors or dotted version vectors; Spanner uses TrueTime (bounded clock uncertainty plus commit wait) when it needs a real-time order.
Design a URL shortener such as TinyURL.
- Requirements: shorten a URL, redirect, optional custom alias and expiry, analytics. Read-heavy (10:1 or more), low-latency redirects, highly available; links must never point to the wrong target.
- Estimates: 100 M new links per day is about 1,200 writes per second and 12,000 reads per second; about 90 TB over five years; 7 base62 characters give 3.5 trillion keys.
- API:
POST /v1/urlsreturns the short URL;GET /{key}returns 301 or 302. - Keys: base62 of a unique 64-bit ID from a Snowflake-style generator or pre-allocated counter ranges (no collisions); hashing plus collision check is the alternative; custom aliases use a conditional insert.
- Data: KV store (key to long URL, owner, expiry), sharded by key and replicated.
- Reads: CDN and Redis cache-aside in front of the KV store; hot links give a high hit ratio.
- Analytics: click events to a stream asynchronously, aggregated offline.
- Trade-offs: 301 reduces load but hides repeat clicks from analytics; sequential IDs are guessable (shuffle bits if that matters); rate limit creation and scan for malicious links.
Design a distributed rate limiter for an API gateway.
- Requirements: per-client limits by rule (for example 100 requests per minute per API key), about 1 ms of overhead, accurate within a few percent, highly available, clear client feedback.
- Algorithm: token bucket (bursts allowed, average enforced) or sliding window counter (accurate with O(1) memory).
- Where: middleware in the gateway, with rules loaded from configuration and cached locally.
- State: Redis keyed by client and rule; a Lua script atomically refills, checks and decrements, and sets an expiry. Shard Redis by key.
- Response: 429 with Retry-After and remaining-quota headers.
- Deep dive: race conditions without atomic scripts; local buckets synced periodically for very high QPS; multi-region limits (per region or asynchronously aggregated); fail open vs fail closed when Redis is down; hot keys for a single abusive client.
Design a chat system such as WhatsApp.
- Requirements: one-to-one and group chat, delivery and read receipts, presence, offline delivery, history across devices, media; messages never lost, ordered within a conversation, end-to-end latency under a few hundred ms.
- Estimates: 500 M daily users x 40 messages is 20 B messages per day, about 230,000 per second average; tens of millions of concurrent connections.
- Connections: WebSocket gateways; a session registry maps user to gateway node with TTL heartbeats.
- Send: gateway to chat service, which assigns a per-conversation sequence number, persists, acknowledges the sender ("sent"), then routes to the recipient's gateway or to push (APNs/FCM) if offline.
- Storage: wide-column store partitioned by conversation_id, clustered by sequence; clients sync with "messages after sequence N". Media in object storage via pre-signed URLs, sent as references.
- Groups: fan-out via a queue for small groups; pull model for very large channels.
- Deep dive: idempotent resends with client message IDs; ordering and multi-device sync; presence throttling; end-to-end encryption with the Signal protocol meaning the server sees only ciphertext; gateway draining during deploys.
Design a news feed such as Twitter/X or Instagram.
- Requirements: post, follow, view a ranked home timeline quickly; eventual consistency is acceptable (seconds of delay).
- Estimates: read-heavy by roughly 40-100 to 1; a celebrity post may need to reach tens of millions of timelines.
- Data: posts store (sharded by post ID), social graph store, timeline cache (Redis lists of post IDs per user, capped to a few hundred), media on a CDN.
- Write path: post saved, event on a queue, fan-out workers push the post ID into followers' timelines (fan-out on write).
- Hybrid: skip fan-out for accounts with huge follower counts; at read time, merge their recent posts into the precomputed timeline (fan-out on read). Skip inactive followers.
- Read path: fetch timeline IDs, hydrate posts from cache, rank (recency, affinity, engagement model), cursor pagination.
- Trade-offs: write amplification vs read latency; storage for precomputed timelines; ranking freshness vs cost.
Design a notification system (push, SMS, email).
- Requirements: send to millions of users across channels, respect preferences and quiet hours, no duplicates, prioritise urgent messages, track delivery.
- Flow: producer services call a notification API with an idempotency key, or emit events; the service validates, checks preferences and rate limits, renders templates, and enqueues to per-channel, per-priority queues.
- Workers: channel workers call provider adapters (APNs, FCM, SMS gateways, email providers) with retries and exponential backoff, failover between providers, and a dead-letter queue.
- Data: device registry (tokens per user, pruned when providers report invalid tokens), preferences, templates, notification log for dedupe and audit.
- Scale: shard by user ID; bulk campaigns paced by a scheduler to respect provider quotas.
- Deep dive: exactly-once effect via dedupe on idempotency key; per-user frequency caps; delivery and open tracking; time-zone-aware scheduling.
Design file storage and sync such as Dropbox or Google Drive.
- Requirements: upload, download, share, sync across devices, version history; files up to many GB; durable; works on flaky networks.
- Chunking: split files into about 4 MB chunks named by content hash; upload only missing chunks (deduplication and delta sync); resumable uploads.
- Services: metadata service (files, folders, versions, chunk lists, ACLs) in a sharded relational database; block storage in an object store; clients upload and download chunks directly via pre-signed URLs.
- Sync: a per-user change journal with increasing versions; devices hold a cursor and long-poll or subscribe for notifications, then pull changes after their cursor.
- Conflicts: optimistic concurrency on file version; on conflict create a "conflicted copy".
- Deep dive: durability through replication or erasure coding, cold storage tiers, garbage collection of unreferenced chunks, sharing permissions, and bandwidth throttling on the client.
Design a video streaming platform such as YouTube.
- Requirements: upload, process, stream at scale globally with smooth playback on varying networks; search and recommendations are separate subsystems.
- Upload: resumable chunked uploads to object storage; an event starts processing.
- Processing: a DAG of jobs on queues: validate, split into segments, transcode into multiple resolutions and codecs in parallel, generate thumbnails, run content moderation, package HLS/DASH manifests.
- Delivery: CDN serves segments; adaptive bitrate players choose quality per segment; popular videos pre-positioned at edges, the long tail served from regional caches or origin.
- Metadata: video info, channels and comments in sharded databases with caches; view counts aggregated via streams.
- Trade-offs: storage and egress cost vs quality (encode popular videos with more renditions and better codecs); processing latency vs cost; eventual consistency for counts.
Design a ride-hailing service such as Uber.
- Requirements: riders request trips, see nearby drivers, get matched quickly; drivers stream location; trips are tracked and billed.
- Estimates: 1 M active drivers updating every 4 seconds is about 250,000 location writes per second.
- Location service: in-memory geospatial index (geohash, quadtree, or S2/H3 cells) sharded by region; only latest position kept hot, history streamed to a log.
- Matching: query the rider's cell and neighbours, rank candidates by ETA, offer to one driver with a timeout, then the next; assignment must be exclusive (conditional update or lock on driver state).
- Trip service: durable state machine (requested, accepted, arriving, in trip, completed, cancelled); payments via a separate idempotent payment service.
- Deep dive: hot cells in city centres (split cells adaptively), surge pricing per cell from supply and demand, ETA from a routing engine, regional isolation for availability.
Design search autocomplete (typeahead).
- Requirements: top 5-10 suggestions per prefix within about 100 ms end to end; popularity-based, optionally personalised; filtered for unsafe terms.
- Data collection: query logs aggregated offline (for example hourly) into prefix frequencies, with decay for freshness.
- Serving structure: a trie with the top-k completions precomputed at every node, so a lookup is O(prefix length); built offline and loaded into memory on serving nodes.
- Scale: shard by prefix range, replicate for reads, cache hot prefixes at the CDN or edge, and client-side debouncing and caching.
- Trade-offs: freshness vs build cost (add a small real-time trending layer), memory vs coverage (only keep prefixes above a frequency threshold).
Design a distributed key-value store.
- Requirements: get and put by key, horizontally scalable, highly available, tunable consistency, durable.
- Partitioning: consistent hashing with virtual nodes.
- Replication: N replicas on successive distinct nodes across zones.
- Consistency: configurable quorums (R + W > N for stronger reads); version vectors to detect concurrent writes; last-writer-wins or client merge for conflicts.
- Failures: gossip-based membership and failure detection, hinted handoff for temporarily down nodes, read repair, and Merkle-tree anti-entropy to reconcile replicas.
- Storage engine: write-ahead log, memtable, SSTables with Bloom filters, compaction, tombstones for deletes.
- Trade-offs: AP with eventual consistency by default vs CP via consensus per shard (Raft groups, like etcd or TiKV) for linearizable operations at higher latency.
Design a distributed cache.
- Requirements: sub-millisecond get and set, horizontal scale to terabytes of RAM, TTLs and eviction, survives node failures (data loss acceptable but should be limited).
- Partitioning: hash slots or consistent hashing; clients or a proxy route by key; slot map updates on resharding.
- Replication: each primary has replicas; automatic failover promotes a replica; asynchronous replication may lose recent writes.
- Node design: hash table in memory, approximate LRU or LFU eviction under a memory cap, lazy plus periodic TTL expiry, event-loop networking.
- Deep dive: hot keys (client-side near cache, key replication), cold-start protection for the database, cache stampede controls, and large values (compress or split).
Design a web crawler.
- Requirements: crawl billions of pages, polite to sites, avoid duplicates, prioritise important and frequently changing pages, extensible for indexing.
- Components: seed URLs, URL frontier with priority queues and per-host queues (politeness delays,
robots.txtcache), DNS cache, fetcher workers, parser and link extractor, URL normaliser and "seen" filter (Bloom filter or hash store), content dedupe (checksums, SimHash for near duplicates), page storage in object storage, index pipeline. - Scale: partition the frontier by host hash so each host is handled by one worker group, enabling politeness without coordination.
- Deep dive: crawler traps (depth limits, URL patterns), recrawl scheduling by change frequency, handling JavaScript-heavy pages with a rendering tier, fault tolerance by checkpointing the frontier.
Design a payment system.
- Requirements: charge customers through external payment providers, never double-charge or lose money, full audit trail, strong consistency, reconcile with providers.
- API:
POST /paymentswith a client-supplied idempotency key; asynchronous status via webhooks and polling. - Flow: payment service records the intent (pending) in a relational database, calls the provider with the same idempotency key, updates status on the response or webhook; retries are safe because both layers dedupe.
- Ledger: double-entry bookkeeping (every movement debits one account and credits another), append-only, so balances are derived and auditable.
- Reliability: outbox pattern for events to downstream systems (orders, notifications); timeouts are "unknown" states resolved by querying the provider; daily reconciliation against provider settlement files flags mismatches.
- Trade-offs: choose CP for the ledger; accept higher latency; separate the synchronous authorisation path from asynchronous settlement; strict security (tokenised card data, PCI scope minimised).
Design a distributed job scheduler (cron at scale).
- Requirements: schedule one-off and recurring jobs, run each at least once near its time, retries, visibility into history, millions of jobs.
- Data: jobs table with schedule, next_run_at, status, owner; index on next_run_at; execution history table.
- Scheduler: partitioned scheduler nodes each own a shard of jobs; they poll for due jobs and claim each atomically (conditional update of status and lease), then enqueue it.
- Execution: worker pools consume from queues, heartbeat while running; if a lease expires the job is retried on another worker; results recorded; next_run_at computed for recurring jobs.
- Deep dive: jobs must be idempotent (at-least-once); avoiding thundering herds at the top of the hour (jitter); priorities and quotas per tenant; job dependencies as a DAG; clock skew (use the database's clock or a single time source).
Design a ticket booking system that handles flash sales.
- Requirements: never oversell, handle huge spikes when sales open, fair, holds that expire if not paid.
- Inventory: seats or ticket counts in a strongly consistent store; reserve with a conditional update (
UPDATE ... SET status='held' WHERE id=? AND status='free') or an atomic decrement in Redis backed by the database. - Holds: a held seat has an expiry (for example 10 minutes); a background job or TTL releases unpaid holds.
- Spike handling: a virtual waiting room (queue with tokens) admits users at the rate the backend can handle; static pages on the CDN; rate limit per user and bot detection.
- Payment: idempotent payment call; on success confirm the booking; on failure release the hold (a saga).
- Trade-offs: strong consistency on inventory vs throughput (shard inventory by event or section, pre-split counts into buckets); fairness vs simplicity.
Design an OTA update system for 100 million Android devices.
- Requirements: deliver signed OS updates safely to a heterogeneous fleet, never brick devices, control rollout, minimise bandwidth and user disruption, observe outcomes.
- Build side: full and delta payloads per source build fingerprint; signed payloads and metadata; stored in a payload store behind a CDN.
- Policy server: devices check in with fingerprint, model, carrier and region (with jitter); the server returns the eligible update according to rollout rules.
- Rollout: staged percentages (1, 10, 50, 100) per cohort, automatic halt when install failures, boot failures or crash rates exceed baseline, and a kill switch via a signed manifest.
- Device: resumable download preferring Wi-Fi and charging; signature and hash verification; install to the inactive A/B slot in the background; reboot into the new slot; if boot does not complete successfully, automatically fall back to the old slot; anti-rollback protection.
- Telemetry: report each stage's outcome to a pipeline that powers rollout dashboards and health gates.
- Trade-offs: delta size vs number of source builds to support; speed of rollout vs blast radius; storage cost of A/B vs virtual A/B snapshots.
Design a telemetry pipeline that collects quality metrics (call drops, crashes) from millions of devices.
- Device: an SDK logs events into a bounded ring buffer, aggregates locally (counts, histograms) instead of shipping raw events, samples high-volume events, scrubs personal data, and honours consent.
- Upload: batched, compressed uploads on Wi-Fi or charging or piggybacked on other traffic; exponential backoff; daily data cap.
- Ingest: regional HTTPS collectors authenticate devices, validate schema versions, and write to a partitioned stream.
- Processing: stream jobs for near-real-time dashboards and alerts; batch jobs into a warehouse partitioned by date, build, model, carrier and region.
- Analysis: compare a new build's metrics against a baseline cohort with confidence intervals to separate real regressions from noise; feed rollout health gates.
- Trade-offs: data freshness vs battery and bandwidth; detail vs privacy; sampling rate vs statistical power.
Design a real-time game leaderboard.
- Requirements: update scores in real time, show the global top 100, show a player's rank and nearby players, millions of players.
- Core structure: Redis sorted set (skip list plus hash):
ZADD,ZREVRANGEfor top N,ZREVRANKfor rank, all O(log n). - Durability: scores persisted in a database; Redis rebuilt from it if lost.
- Scale: one sorted set handles tens of millions of members; beyond that, shard by score range (rank = rank within shard + counts of higher shards) or compute approximate ranks with score histograms for players outside the top.
- Extras: per-region and per-period leaderboards (daily keys with expiry); friend leaderboards computed on demand from the friend list.
Design a metrics and monitoring system (like Prometheus plus Grafana at scale).
- Requirements: ingest millions of samples per second from services, query recent data fast for dashboards and alerts, retain history cheaply.
- Collection: agents scrape or receive metrics, pre-aggregate, and forward to ingesters; labels define series, and cardinality must be controlled.
- Storage: time-series database with per-series compressed chunks (delta-of-delta timestamps, XOR-compressed values), sharded by series hash, replicated.
- Retention: raw data for days, downsampled rollups (1 minute, 1 hour) for months, in object storage.
- Alerting: a rule engine evaluates queries periodically, deduplicates and routes alerts with silencing and escalation.
- Deep dive: high-cardinality labels (user IDs) explode series counts; the monitoring system must be more reliable than what it monitors, with a separate failure domain.
Design a collaborative document editor such as Google Docs.
- Requirements: multiple users edit simultaneously with low latency, all converge to the same document, offline edits merge later, history and permissions.
- Concurrency control: operational transformation (a central server transforms concurrent operations against each other) or CRDTs (operations designed to commute, allowing peer-to-peer and offline merging).
- Architecture: clients apply edits locally immediately (optimistic), send operations over WebSocket to a document session server that owns the document (routed by document ID), which orders, transforms and broadcasts them.
- Storage: an operation log plus periodic snapshots for fast loading and version history.
- Extras: presence and cursors as ephemeral state; permission checks per operation; session server failover by replaying the log.
Design an ad-click aggregation system.
- Requirements: count clicks per ad per minute for billing and dashboards; billions of events per day; accurate (money is involved); results within minutes; handle late and duplicate events.
- Ingest: click events with unique IDs to a partitioned log (partition by ad ID).
- Processing: a stream processor (Flink-style) with event-time tumbling windows, watermarks for late data, and exactly-once state via checkpoints; dedupe by click ID.
- Output: aggregated counts to an OLAP store for queries; raw events archived to object storage.
- Correctness: a daily batch job recomputes counts from raw events and reconciles with the streaming results (lambda-style check); fraud filtering before billing.
- Trade-offs: latency vs completeness (how long to wait for late events); hot ads (split partitions with key salting).
Scenario & debugging
Your primary database is at 90 percent CPU and read traffic keeps growing. What do you do, in order?
- Find the expensive queries (slow query log, query plans) and fix them: missing indexes, N+1 queries, unbounded scans.
- Add caching for hot, read-mostly data (cache-aside with TTLs) and a CDN for cacheable responses.
- Add read replicas and route read-only queries to them, handling replication lag for read-your-writes flows.
- Move analytics and reporting queries to a warehouse fed by change data capture.
- Scale up the instance as a short-term relief while doing the above.
- If writes are also the bottleneck, split by function (separate databases per domain) and eventually shard by a well-chosen key.
p99 latency tripled right after a deploy. How do you investigate?
First mitigate: if the timing correlates with the deploy, roll back or disable the feature flag; restore service before root-causing. Then compare the new and old versions: traces for the slow requests (which span grew), metrics for CPU, garbage collection, thread and connection pool saturation, and dependency latency. Common culprits: a new synchronous call on the hot path, a query without an index, a cache key change causing misses, a reduced pool size, or logging at high volume. Add a canary stage with automatic latency comparison so the next regression is caught at 1 percent of traffic.
A celebrity's account overloads one shard every time they post. How do you fix it?
This is a hot-key problem. Serve reads of the celebrity's profile and posts from caches replicated across many nodes (and a local in-process cache with a short TTL). For timelines, stop fanning out their posts on write; merge them at read time. For writes such as likes and counters, split the counter into many sub-counters (key salting) and sum them asynchronously. If one tenant is consistently hot, move it to a dedicated shard via a directory-based mapping.
The Redis cache cluster restarted and the database immediately fell over. What went wrong and how do you prevent it?
A cold cache means every request misses and hits the database at once (a cache avalanche), far beyond its capacity. Prevention: replicas and persistence so the cache survives node restarts; rolling restarts one shard at a time; warming the cache from a snapshot or with the hottest keys before taking traffic; request coalescing so only one request per key goes to the database; circuit breakers and load shedding in front of the database; and admission control that ramps traffic up gradually.
Customers report being charged twice. How do you investigate and fix it?
Look for the retry path: a client or gateway retrying after a timeout, a queue redelivering a message, or a double-click in the UI. Confirm in logs that two provider charges share the same order but different requests. Fix: require an idempotency key per payment attempt, store it with a unique constraint before calling the provider, pass the same key to the provider (most support it), and return the stored result for duplicates. Make consumers dedupe by message ID. Then refund the affected customers and add reconciliation that alerts on duplicate charges.
Consumer lag on a Kafka topic keeps growing. What do you check?
Check whether producers increased their rate (a spike or a bug), whether consumers slowed down (a slow downstream database, expensive processing, garbage collection), and whether some partitions lag more than others (a hot partition key, or a consumer stuck on a poison message and retrying). Remedies: scale consumers up to the partition count (and add partitions for future scale), batch downstream writes, move poison messages to a dead-letter topic, fix skewed keys, and apply backpressure or load shedding upstream. Alert on lag growth, not just absolute lag.
How would you migrate a large single database to a sharded cluster with no downtime?
- Introduce a data-access layer that knows the shard key and routing.
- Set up the new shards and start dual writes (or stream changes via change data capture) from the old database to the new cluster.
- Backfill historical data in batches, then verify with checksums and sampled comparisons.
- Shadow-read from the new cluster and compare results with the old.
- Gradually switch reads, then writes, per tenant or percentage behind feature flags.
- Keep the old database in sync for a rollback window, then decommission.
The same expand, migrate, contract approach applies to schema changes.
An entire cloud region goes down. How should your system respond?
It depends on the multi-region design chosen in advance. Active-passive: health checks fail, DNS or global load balancing shifts traffic to the standby region, and a replica there is promoted; RPO equals the replication lag, RTO the failover time. Active-active: traffic is already served from several regions, so the remaining ones absorb load (they must have headroom), with data replicated asynchronously and conflicts resolved. Either way, stateless tiers must be pre-deployed, capacity reserved, and failover rehearsed regularly (game days); otherwise the plan fails when needed.
Users say their profile edits "disappear" and then reappear a few seconds later. What is happening?
Classic replication lag or cache staleness: the write went to the leader, but the next read hit a lagging follower or a cache entry that was not invalidated. Fixes: read-your-writes by routing the user's reads to the leader for a short period after writing (or waiting for the follower to reach the write's log position); invalidate the cache after the database commit; pin a user's session to one replica for monotonic reads; and, on the client, optimistically show the edited value.
A dependency slows down and your whole service falls over, even endpoints that do not use it. Why, and how do you prevent it?
Requests to the slow dependency held threads or connections for a long time, exhausting a shared pool so every endpoint queued (a cascading failure), and retries multiplied the load. Prevent with tight timeouts, a circuit breaker around the dependency, bulkheads (a separate, bounded pool per dependency), a retry budget with backoff and jitter, load shedding when queues grow, and fallbacks such as cached or default responses.
At 10 percent of an OTA rollout, boot-failure telemetry rises above baseline. What do you do?
The rollout health gate should already have paused the rollout automatically; if not, halt it immediately via the policy server or kill switch so no new devices receive the build. Affected devices with A/B updates should fall back to the previous slot automatically; confirm from telemetry that they do. Slice the failures by model, carrier, region and source build to find the common factor (for example a delta payload for one source build, or one hardware variant). Fix, rebuild, and restart the rollout from 1 percent for the affected cohort. Afterwards, tighten the gate thresholds and add the failing configuration to pre-release testing.
After shipping a new analytics SDK, users complain about battery drain. How do you diagnose and fix it?
Compare battery and wakeup metrics between devices with and without the SDK version. Typical causes: frequent timers or alarms waking the CPU, uploads on every event instead of batching, holding wakelocks without timeouts, keeping the radio active with small periodic transfers, or retrying aggressively when offline. Fix by batching events in a bounded buffer, scheduling uploads with the OS job scheduler on Wi-Fi and charging, coalescing with other network activity, using backoff with a cap, and adding a remote kill switch or config to reduce sampling in the field.
Design how per-carrier settings get delivered to millions of phones, overridable by SIM and region.
- Layering: platform defaults, then carrier-specific values keyed by the SIM's network code, then server-pushed overrides, then a runtime cache that apps read through one API.
- Delivery: signed, versioned config bundles fetched from a CDN on check-in or pushed via a notification that triggers a fetch.
- Invalidation: recompute on SIM swap, roaming changes and OS updates; versions let devices detect staleness.
- Safety: staged rollout by percentage and cohort, validation on device with fallback to the last known good config, a kill switch, and telemetry on apply success and key quality metrics.
On a dual-SIM phone, the radio HAL keeps crashing and framework calls hang. How would you make the system resilient?
Register death notifications on each HAL service; when the HAL dies, immediately fail every in-flight request with a "radio not available" error so callers do not hang. Put timeouts on every request (and on any wakelock held while waiting) so a lost response cannot deadlock callers or pin the CPU awake. On HAL restart, re-initialise idempotently with bounded exponential backoff, and never crash the long-lived framework process on a transient error. Add backpressure by capping outstanding requests and coalescing duplicate polls. Emit metrics on HAL restarts and request timeouts so fleet dashboards catch regressions.
Leadership asks you to cut the infrastructure cost of your service in half. Where do you look?
Measure first: break down cost by component (compute, storage, database, CDN egress, cross-zone traffic, logging). Then: right-size over-provisioned instances and use autoscaling; use spot or preemptible capacity for batch work; raise cache hit ratios to shrink database size; move cold data to cheaper storage tiers and set retention policies; compress payloads and cut egress with better CDN caching; reduce log and metric volume (sampling, lower cardinality); and remove idle environments. Name the reliability or latency trade-off for each change.
The interviewer says "traffic just grew 10x". How do you respond?
Walk the request path and ask where each component breaks. Stateless services: add instances behind the load balancer and confirm autoscaling limits. Cache: more nodes, check hot keys. Database: reads go to more replicas and cache; writes may now exceed one primary, so shard by the key identified earlier. Queues: add partitions and consumers. Fan-out-heavy paths (feeds, notifications) need hybrid strategies. Check limits you do not own: third-party API quotas, connection limits, cross-region bandwidth. Finally, revisit estimates and cost, since 10x traffic is usually 10x bill unless efficiency improves.