Core Computer Science

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.

~115 min read 0 interview questions
In 30 seconds
  • 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).

Analogy

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

AvailabilityDowntime per yearDowntime per month
99% (two nines)about 3.65 daysabout 7.3 hours
99.9% (three nines)about 8.8 hoursabout 44 minutes
99.99% (four nines)about 53 minutesabout 4.4 minutes
99.999% (five nines)about 5.3 minutesabout 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.
Common pitfall Claiming five nines while every request passes through one database in one data centre. Availability is limited by the weakest serial dependency; show where the redundancy is.
Interview angle Interviewers check whether you turn "it should be fast and reliable" into numbers: a p99 latency target, an availability SLO, peak QPS. Strong candidates explain how statelessness allows horizontal scaling, quote percentiles instead of averages, and do quick availability arithmetic for serial and redundant components.

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.

Analogy

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)

ModelGuaranteeTypical use
Linearizable (strong)Every operation appears to take effect at one instant between its start and end; everyone sees the same latest valueLocks, leader election, balances, unique usernames
SequentialAll nodes see operations in the same order, consistent with each client's own order, but not necessarily real timeSome replicated logs
CausalOperations that are causally related (a reply after a post) are seen in order by everyone; unrelated ones may differComments and replies, collaborative apps
Read-your-writesA client always sees its own writesProfile edits, settings
Monotonic readsA client never sees data go backwards in timePinning a user to one replica
EventualIf writes stop, all replicas converge eventuallyLikes, 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.

Tip Choose consistency per feature, not per system. In one e-commerce app, payments and inventory need strong consistency, while product reviews, recommendations and view counts can be eventually consistent.
Interview angle "Explain CAP with a real choice" is near-universal. Strong answers note that CAP only bites during partitions, that the C in CAP means linearizability (not the C in ACID), bring in PACELC for the normal-case latency trade-off, and pick a model per data type with a reason.

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.

Analogy

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

QuantityValue
Seconds per day86,400 (round to 10^5)
1 million requests per dayabout 12 per second on average
Peak trafficabout 2-3x average (more for spiky events)
Powers of tenthousand 10^3 (KB), million 10^6 (MB), billion 10^9 (GB), trillion 10^12 (TB)
Common sizesint 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 serverroughly 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)

OperationApproximate time
L1 cache referenceabout 1 ns
Main memory referenceabout 100 ns
Read 1 MB sequentially from memoryabout 10-50 microseconds (classic 2012 number was ~250 μs)
SSD random readabout 16-100 microseconds
Round trip within one data centreabout 0.5 ms
Read 1 MB sequentially from SSDabout 0.2-1 ms
Hard disk seekabout 10 ms
Send 1 MB over a 1 Gbps linkabout 10 ms
Round trip California to Netherlandsabout 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.

Tip Say every assumption out loud and write it down. Interviewers grade the reasoning and whether the numbers change your design, not arithmetic precision. Spend three to five minutes here, not fifteen.
Interview angle Expect "how much storage will this need in five years?" or "how many servers?". Interviewers look for sensible assumptions, correct orders of magnitude, separate read and write QPS, peak vs average, and, most importantly, a design decision that follows from each number.

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).

Analogy

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

AspectLayer 4 (transport)Layer 7 (application)
Decides usingIP addresses and portsHTTP method, path, headers, cookies
SpeedVery fast, protocol agnosticSlightly more CPU (parses requests)
FeaturesConnection-level balancingPath-based routing, TLS termination, sticky sessions, retries, compression, header rewriting
ExamplesAWS NLB, LVSAWS 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.js so you never need to purge; purge APIs are slow and eventually consistent.
  • Signed URLs protect private content on the CDN with expiring tokens.
Common pitfall Relying on sticky sessions for correctness. If the pinned server dies, the user's session dies with it. Keep servers stateless and use stickiness only as a cache-locality optimisation.
Interview angle Interviewers ask L4 vs L7, which balancing algorithm for long-lived WebSocket connections (least connections, plus graceful draining on deploy), how to avoid the load balancer being a single point of failure, and when a CDN helps (static and cacheable content, global users) and when it does not (personalised, write-heavy traffic).

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).

Analogy

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

StrategyHow it worksProsCons
Cache-aside (lazy loading)App reads cache; on miss, reads DB and fills cacheMost common; cache holds only requested data; cache failure is not fatalFirst read is slow; stale until TTL or invalidation
Read-throughCache library loads from DB on miss itselfSimpler app codeSame staleness issues; cache must know the DB
Write-throughWrite to cache and DB synchronouslyCache always freshSlower writes; caches data that may never be read
Write-back (write-behind)Write to cache; flush to DB asynchronouslyVery fast writes; batchingData loss if cache node dies before flush
Write-aroundWrite to DB only; cache filled on later readsAvoids polluting cache with write-once dataRecent writes miss the cache
Refresh-aheadRefresh hot keys before they expireAvoids miss latency for hot dataWasted 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

ProblemWhat happensFixes
Stampede (thundering herd)A hot key expires and thousands of requests hit the DB at onceRequest coalescing or a per-key lock (only one rebuilds), stale-while-revalidate, TTL jitter, early probabilistic refresh
PenetrationRequests for keys that do not exist always miss and hit the DBCache negative results briefly; Bloom filter of existing keys
AvalancheMany keys expire together, or the cache cluster restarts coldRandomised TTLs, warm-up, replicas, circuit breaker to protect the DB
Hot keyOne key (a celebrity profile) overloads one cache shardReplicate the key across shards with suffixes, local in-process cache, CDN
Staleness and inconsistencyCache and DB disagreeTTL, 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
Common pitfall Adding a cache without saying how it is invalidated, what happens on a cold start, or what the hit ratio needs to be. A cache with a 20 percent hit rate adds a network hop and complexity for little gain.
Interview angle Expect "which caching strategy and why?", "how do you keep the cache consistent with the database?", "what happens when a hot key expires?" and "how do you implement LRU in O(1)?". Strong answers quantify cache size, name the eviction policy, handle stampedes and cold starts, and state how much staleness the product can tolerate.

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.

Analogy

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

TypeExamplesStrengthsUse for
Relational (SQL)PostgreSQL, MySQLSchemas, joins, ACID transactions, mature toolingPayments, orders, users, anything relational with integrity constraints
Key-valueDynamoDB, Redis, RocksDBSimple, very fast lookups by key; easy horizontal scaleSessions, URL mappings, carts, feature flags
DocumentMongoDB, CouchbaseFlexible JSON-like records, nested dataCatalogues, content, user profiles
Wide-columnCassandra, HBase, BigtableMassive write throughput, partition key plus clustering orderTime series, messages, activity logs
GraphNeo4j, NeptuneFast multi-hop relationship queriesSocial graphs, fraud rings, recommendations
SearchElasticsearch, OpenSearchInverted index, full-text and faceted searchSearch boxes, log analytics
Time-seriesInfluxDB, Prometheus, TimescaleDBCompression and downsampling by timeMetrics, telemetry
Blob / object storageS3, GCSCheap, durable storage of large filesImages, video, backups, logs
Columnar warehouseBigQuery, Redshift, ClickHouseFast 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

TopologyHow it worksTrade-offs
Leader-follower (primary-replica)All writes go to the leader, which streams changes to followers; reads can go to followersSimple and common. Async replication means lag, so stale reads and possible loss of recent writes on failover
Multi-leaderSeveral leaders accept writes (for example one per region) and replicate to each otherLocal writes in each region, but write conflicts must be resolved
Leaderless (Dynamo style)Clients write to and read from several replicas using quorumsHigh 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.

SchemeHowProsCons
RangeKey ranges per shard (A-F, G-M, ...)Efficient range scansHot spots on sequential keys such as timestamps
Hashhash(key) mod N or consistent hashingEven spreadRange queries hit all shards; mod N remaps nearly everything when N changes
Directory / lookupA service maps key to shardFlexible, easy to move tenantsLookup service is an extra dependency
GeographicBy user regionData locality, residency complianceUneven 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
Tip Scale in this order and say so: optimise queries and add indexes, add caching, add read replicas, split by function (separate databases per service), and only then shard. Each step adds complexity, so do not shard a database that fits on one machine.
Interview angle Interviewers want the database choice justified by access patterns ("messages are written once, read by conversation in time order, so a wide-column store partitioned by conversation_id and clustered by timestamp"). Expect follow-ups on shard key choice and hot shards, replication lag and read-your-writes, what an index costs, and how consistent hashing limits data movement.

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.

Analogy

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

SemanticMeaningHow
At-most-onceMay lose messages, never duplicatesAcknowledge before processing
At-least-onceNever loses, may duplicateAcknowledge after processing; the common default
Exactly-once (effectively)Each message affects state onceAt-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_id or order_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.
Common pitfall Adding a queue and assuming it gives exactly-once processing and global ordering. It gives neither by default: you get at-least-once delivery and per-partition order, so design idempotent consumers and choose partition keys deliberately.
Interview angle Interviewers ask why you need a queue here, Kafka vs RabbitMQ, how you guarantee ordering, what happens when a consumer crashes mid-message, how to avoid double-charging, and the dual-write problem. Naming idempotency keys, DLQs, the outbox pattern and consumer lag signals real production experience.

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.

Analogy

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

StyleStrengthsWeaknessesTypical use
REST over HTTP/JSONSimple, cacheable, universal toolingOver- or under-fetching; loose contractsPublic APIs, CRUD resources
gRPC (HTTP/2 + Protocol Buffers)Fast binary encoding, strict schemas, streaming, generated clientsHarder to use from browsers; less human-readableInternal service-to-service calls
GraphQLClient picks exactly the fields it needs; one round tripCaching and cost control are harder; N+1 query riskMobile and web clients aggregating many resources
WebSocket / server-sent eventsServer push, low-latency bidirectional updatesStateful connections to manage and scaleChat, live dashboards, presence
WebhooksServer-to-server push on eventsReceiver must be reachable; retries and signatures neededPayment callbacks, integrations

Good API design checklist

  • Resource-oriented URLs and correct methods: GET /users/{id}, POST /orders, PATCH for 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

AlgorithmHow it worksProsCons
Token bucketBucket of capacity B refills at R tokens per second; each request takes a tokenAllows short bursts up to B while enforcing average rate R; tiny stateBurst may still hurt a fragile backend
Leaky bucketRequests enter a queue that drains at a fixed rateSmooth, constant output rateBursts are delayed or dropped; adds latency
Fixed window counterCount requests per clock window (per minute)SimplestUp to 2x the limit across a window boundary
Sliding window logStore a timestamp per request; count those within the last windowExactMemory proportional to request volume
Sliding window counterWeighted blend of the current and previous fixed windowsClose to exact with O(1) memoryApproximate
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.

Interview angle "Design a rate limiter" is a classic. Interviewers want an algorithm choice with a reason (token bucket for burst tolerance), where it runs (gateway vs service vs client), how state is shared across nodes without races, how it fails (fail open vs fail closed if Redis is down), and what the client sees. For API design they check resource modelling, idempotency, pagination and versioning.

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.

Analogy

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

PatternWhat it doesWatch out for
TimeoutsBound how long you wait for any remote callWithout them, slow dependencies exhaust threads and cascade
Retries with exponential backoff and jitterRecover from transient failuresOnly retry idempotent operations; cap attempts; use retry budgets to avoid retry storms
Circuit breakerAfter repeated failures, fail fast for a while, then probe (half-open)Needs a sensible fallback
BulkheadSeparate thread or connection pools per dependencyOne slow dependency cannot starve the others
Load sheddingReject low-priority work when overloadedBetter to serve 80 percent well than 100 percent badly
Graceful degradationServe reduced functionality (cached or default data)Decide in advance which features are optional
Health checks and auto-healingReplace unhealthy instances automaticallyDistinguish liveness (restart me) from readiness (do not send traffic yet)
Redundancy across zones and regionsSurvive data-centre failuresActive-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.
Common pitfall Retrying at every layer. If the client, gateway and service each retry three times, one failing request becomes 27 calls to the struggling dependency. Retry at one layer, with a budget.
Interview angle After the happy path, interviewers ask "what happens when X fails?" for each box. Strong candidates proactively walk through failure modes (node, zone, dependency, overload, bad deploy), name the containment pattern for each, and close with how they would observe and roll out the system safely.

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.

Analogy

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.

  1. 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.
  2. 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.
  3. 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.
  4. Data model (3-5 min) Entities, relationships, access patterns, SQL vs NoSQL with a reason, primary keys and partition keys.
  5. 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.
  6. 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.
  7. 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.
  8. 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.

Tip Keep a visible checklist in the corner of the whiteboard (requirements, estimates, API, data, HLD, deep dive, trade-offs, ops) and tick items off. It keeps you on time and shows the interviewer you are driving.
Common pitfall Naming technologies instead of reasoning ("I'll use Kafka and Cassandra"). Always attach the requirement that justifies the choice ("writes are append-only and need to scale to 50,000 per second, so a wide-column store partitioned by conversation").
Interview angle Interviewers score requirements gathering, a coherent high-level design, depth in at least one area, trade-off reasoning, and communication. They often steer the deep dive deliberately to see whether you can adapt; follow their lead rather than defending your plan.

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.

Analogy

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/urls returns 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-After and X-RateLimit-Remaining headers.
  • 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.
Interview angle For these "infrastructure" designs, interviewers go deep on correctness under concurrency and failure: key collisions in the shortener, race conditions in the rate limiter, conflict resolution in the KV store, duplicate notifications, and double-running cron jobs. Idempotency and atomic claims are the recurring answers.

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.

Analogy

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_id and 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

PromptSignature problem
Payment systemExactly-once effect with idempotency keys, double-entry ledger, reconciliation with the payment provider, strong consistency
Ticket booking / flash saleContention on limited inventory: holds with expiry, conditional updates, queueing users, preventing overselling
LeaderboardSorted sets (skip lists) for rank queries; sharding by score range or approximate global ranks
Metrics and logging platformHigh-volume ingestion, time-series storage, downsampling, retention tiers
Search engineInverted index, sharding by document, scatter-gather queries, ranking
Collaborative editorOperational transforms or CRDTs for concurrent edits, presence, real-time sync
Ad click aggregationStream processing with windowing, exactly-once counting, late events, reconciliation
Interview angle Each product design has one signature trade-off the interviewer is waiting for: celebrity fan-out in feeds, persist-before-ack and ordering in chat, chunking and conflict handling in file sync, the transcoding pipeline and CDN cost in video, and geospatial indexing plus exclusive assignment in ride-hailing. Get to it early and go deep.

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.

Analogy

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

ConstraintDesign response
Battery and powerBatch 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 limitsThrottle background work, schedule heavy jobs (on-device ML, indexing) when idle and cool
Memory and CPUBounded buffers and caches, streaming instead of loading whole files, low-memory-killer awareness, avoid work on the UI thread
Intermittent, metered networksOffline-first local store with a sync queue, retries with backoff and jitter, resumable transfers, compression, respect metered data
Storage wear and spaceRing buffers and log rotation, quotas, avoid write amplification on flash
Privacy and securityMinimise and aggregate data on device, scrub personal data before upload, user consent, signed and verified code and configuration, hardware-backed keys
Heterogeneous fleetMany OS versions, chipsets, carriers and regions; version every API and config; target rollouts by device fingerprint
Cannot patch instantlyFeature 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:

PatternWhere it appears on the deviceGeneral system-design equivalent
Layering with a stable HALApp, framework, RIL, HAL, vendor modemVersioned service interfaces that let layers evolve independently
Async request/response with correlation IDsRIL request serial numbers matched to HAL responsesRequest IDs in any asynchronous RPC or messaging system
Publish/subscribeRegistrant lists, telephony registry callbacksEvent bus fan-out decoupling producers and consumers
Single-threaded event loopHandler and Looper per componentActor model: serialise state changes without locks
Explicit state machinesService-state, call and data-connection trackersDeterministic handling of complex transitions (trip or order state machines)
Death detection and recoveryHAL death recipients, flushing in-flight requests, watchdog restartHealth checks, supervisors and circuit breakers
Backoff and throttlingNetwork-mandated retry timers, data retry managerExponential backoff with jitter, congestion control
Caching with invalidationCarrier configuration, cached service state, SIM recordsCache with invalidation on a change event
Rate limiting and coalescingThrottled signal-strength and unsolicited indicationsDebouncing and protecting consumers from event storms
Resource arbitrationWhich SIM carries data on a dual-SIM phone, shared modem radioFair and safe scheduling of a scarce shared resource
Power-aware designWakelocks with timeouts, hardware offload, radio sleep cyclesBatching 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.
Tip When given a generic system-design question, anchor it in device experience where you can. Concrete examples of backoff, idempotency, failure detection, resource arbitration and power-aware batching on real devices are more convincing than the textbook versions.
Interview angle Mobile and embedded interviewers probe what happens when the device is offline, low on battery, rebooted mid-operation, or running an old version; how you roll out safely to a heterogeneous fleet; and how you avoid bricking devices. Mentioning A/B partitions with automatic rollback, staged rollouts gated on telemetry, signed configuration, and power-aware batching shows you understand the domain.

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.

Analogy

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

LevelExpected
Mid-levelA working high-level design with the right building blocks; answers the interviewer's probes correctly
SeniorDrives the interview; clarifies requirements and quantifies; goes deep on one or two areas; handles failure modes; justifies technology choices
StaffIdentifies 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
PrincipalFrames 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-offOne sideOther side
Consistency vs availability and latencyCorrectness for money, inventory, identityUptime and speed for feeds, counts, presence
Latency vs throughputInteractive requests, small batchesBatching, compression, pipelines
SQL vs NoSQLTransactions, joins, constraintsHorizontal write scale, flexible schema
Push vs pullLow read latency, write amplificationCheap writes, read-time work
Sync vs asyncSimple, immediate resultDecoupled, resilient, eventually consistent
Normalise vs denormaliseOne source of truth, cheaper writesFast reads, duplicated data to keep in sync
Build vs buy (managed service)Control, custom needsSpeed, less operational burden, vendor lock-in
Monolith vs microservicesSimplicity, fast iteration for small teamsIndependent scaling and deployment for large orgs
Accuracy vs costExact counts, sliding logsApproximate 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.
Common pitfall Over-engineering from the start: sharding, multi-region active-active and microservices for a system with 1,000 users. Senior interviewers reward starting simple, showing where it breaks, and evolving it with numbers.
Interview angle Senior loops ask "why not the alternative?", "what would you do differently with half the budget?", "how would you migrate the existing system without downtime?" and "what keeps you up at night in this design?". Treat the interviewer as a colleague in a design review: invite pushback, admit uncertainty, and reason from requirements.

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.

Analogy

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

  1. 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.
  2. 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.
  3. 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. add sets k bits; maybe_contains is 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.
Interview angle "How does Raft actually work?" should get you through election, majority commit, and "why not put the whole database in etcd?". Bloom filters come up for crawlers, caches and LSM reads; say "false positives, no false negatives, no deletes" in one breath.

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?
LevelPreventsStill allows
Read uncommittedAlmost nothingDirty reads
Read committedDirty readsNon-repeatable reads, phantoms, lost updates
Repeatable read / snapshotNon-repeatable reads (and phantoms in some engines)Write skew
SerializableAll anomaliesNothing, 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/urls returns 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.txt cache), 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 /payments with 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, ZREVRANGE for top N, ZREVRANK for 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?
  1. Find the expensive queries (slow query log, query plans) and fix them: missing indexes, N+1 queries, unbounded scans.
  2. Add caching for hot, read-mostly data (cache-aside with TTLs) and a CDN for cacheable responses.
  3. Add read replicas and route read-only queries to them, handling replication lag for read-your-writes flows.
  4. Move analytics and reporting queries to a warehouse fed by change data capture.
  5. Scale up the instance as a short-term relief while doing the above.
  6. 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?
  1. Introduce a data-access layer that knows the shard key and routing.
  2. Set up the new shards and start dual writes (or stream changes via change data capture) from the old database to the new cluster.
  3. Backfill historical data in batches, then verify with checksums and sampled comparisons.
  4. Shadow-read from the new cluster and compare results with the old.
  5. Gradually switch reads, then writes, per tenant or percentage behind feature flags.
  6. 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.