Every system is the same five boxes.
The interview is about the wires.
Clients, load balancers, stateless services, caches, and durable stores. Twelve topologies cover almost every question you will ever be asked. This manual teaches you the boxes first, then the wires, then twenty case studies, then the object-level design rounds — with runnable Java throughout.
How to use this manual
There are exactly two kinds of design interview, and confusing them is the most common way strong candidates lose rounds.
| HLD (High-Level Design) | LLD (Low-Level Design / Machine Coding) | |
|---|---|---|
| Question sounds like | "Design Instagram." "Design a rate limiter." | "Design a parking lot." "Design Splitwise, write the classes." |
| Output | A boxes-and-arrows diagram + numbers + trade-off discussion | Class diagram + compiling, runnable Java |
| Unit of thought | Services, storage engines, network hops | Classes, interfaces, methods, threads |
| Graded on | Scoping, estimation, bottleneck-finding, trade-offs | Extensibility, SOLID, correct concurrency, clean naming |
| Duration | 45–60 min, whiteboard/Excalidraw | 60–120 min, IDE, sometimes take-home |
| Who asks it | FAANG, most product companies, senior+ roles | Indian product cos (Flipkart, Swiggy, Uber India, Atlassian, Zomato, PhonePe), SDE-1/2 heavy |
Part I (F) is the vocabulary — you cannot design what you cannot name. Part II (H) is twenty worked HLD case studies. Part III (L) is object-oriented design with complete Java. Part IV (X) is interview craft: mock transcripts, scripts, drills, and a study plan.
Suggested path if you have four weeks
- Week 1 — F1 → F10. Do the estimation drills until QPS math is automatic.
- Week 2 — F11 → F19, then H1 → H6. One case study per day, out loud, on paper, timed to 45 minutes.
- Week 3 — L1 → L8 (patterns + concurrency), then L9 → L14 typed from scratch without looking.
- Week 4 — H7 → H20, X1 → X8, and re-do the two case studies you did worst on.
Never read a case study passively. Read only the "Requirements" block, close the page, spend 20 minutes designing it yourself on paper, then read the rest and diff your answer against it. Passive reading produces recognition; diffing produces recall.
The HLD interview framework
The question "Design YouTube" is deliberately unanswerable as stated. Your first job is not to answer it — it is to convert an ambiguous prompt into a bounded engineering problem. Interviewers score that conversion heavily. Here is the four-phase structure, with the minutes you should spend in a 45-minute round.
Phase 1 — Scope & requirements (5–8 min)
Ask questions until you can write down a numbered list. Three categories, always in this order:
- Functional requirements — what a user can do. Keep it to 3–5 verbs. "Upload a video, watch a video, search videos." Explicitly park the rest: "I'll treat comments, monetisation, and recommendations as out of scope unless you want them."
- Non-functional requirements — the properties that decide the architecture. Availability target, latency target, consistency requirement, durability, read:write ratio, scale.
- Scale numbers — DAU, requests/sec, object sizes, retention. Ask for them; if the interviewer says "you tell me," pick round numbers and say them out loud.
"Before I design anything — is this read-heavy or write-heavy, and are we optimising for availability or for consistency?" Every subsequent decision follows from those two answers, and asking shows you know that.
Phase 2 — Estimation & API/data contract (5–8 min)
Do the back-of-envelope math (F2). Then write the API surface — 3–6 endpoints is plenty — and the core entities with their key fields. Doing this before drawing boxes prevents the classic failure of designing infrastructure that cannot express the product.
?page=500 forces the database to skip 10,000 rows on a moving dataset and produces duplicates when items are inserted. Cursor pagination (WHERE (created_at, id) < (:ts, :id) ORDER BY created_at DESC, id DESC LIMIT 20) is O(limit) and stable. Saying this unprompted reads as production experience.
Phase 3 — High-level design, then deep dives (20–25 min)
Draw the simplest thing that satisfies the requirements — usually client → LB → service → DB. Walk one read path and one write path end to end, out loud. Only then start scaling it, and scale it because a number forces you to, not because caching is a reflex.
The interviewer will then pick one or two components and say "go deeper." Common deep dives: the sharding key, the cache invalidation strategy, how the queue guarantees delivery, how you handle a hot partition, what happens when a node dies mid-write.
Phase 4 — Bottlenecks, failure, wrap-up (5 min)
Volunteer the weaknesses before you are asked. "The single point of failure here is X; I'd fix it with Y. The thing I'd monitor first is Z. If traffic grew 10×, the first component to break is the metadata DB, and the fix is to shard on userId."
What the interviewer is actually filling in on the scorecard
| Signal | Strong | Weak |
|---|---|---|
| Requirement handling | Bounds the problem, states assumptions aloud | Starts drawing Kafka in minute two |
| Estimation | Rough numbers that drive decisions | Precise arithmetic that changes nothing |
| Trade-offs | "I'd choose A over B because our read:write is 100:1" | Lists technologies without comparison |
| Depth | Can explain how the chosen DB actually stores data | "We'll use Cassandra because it scales" |
| Communication | Drives the conversation, checks in, adapts | Monologue, ignores hints |
When an interviewer says "hmm, interesting — what happens if two users do that at the same time?", they are not curious. They have found a bug in your design. Stop, engage with it fully, and change the design. Candidates who say "yes, good point" and continue with the same design fail.
Practice — "Design Twitter." What are your first five questions?framework
1. Which features are in scope — post a tweet, home timeline, follow, search? I'd propose the first three and drop search unless you want it.
2. What scale — DAU, average follower count, and is there a celebrity tail? (This single answer determines fan-out-on-write vs fan-out-on-read.)
3. How stale can the home timeline be — is 5 seconds acceptable? (Buys eventual consistency.)
4. Read:write ratio, so I know what to optimise. (Twitter is roughly 1000:1 read-heavy.)
5. Do we need media (images/video), or text only? (Media changes the storage story entirely — blob store + CDN.)
Why these five: each one, if answered differently, changes the architecture. Questions whose answers do not change your design are wasted minutes and the interviewer notices.
Practice — You are 30 minutes in and the interviewer has said nothing for 10 minutes. What do you do?behaviour
Stop and check in: "I've covered the write path and storage. I could either go deeper into how I'd shard the metadata store, or step back and cover the read path and caching — which is more useful to you?" Silence usually means you are in a region they do not care about. Explicitly offering a fork lets them steer without having to interrupt you, and it demonstrates the collaborative behaviour they are grading.
Back-of-the-envelope estimation
The point of estimation is not arithmetic. It is to answer one question: does this fit on one machine? If yes, do not distribute it. If no, the numbers tell you which dimension broke first — CPU, memory, disk, or network — and that dimension dictates the architecture.
Numbers you must know cold
| Operation | Time | Mental model |
|---|---|---|
| L1 cache reference | ~1 ns | free |
| Branch mispredict | ~3 ns | free |
| L2 cache reference | ~4 ns | free |
| Mutex lock/unlock (uncontended) | ~17 ns | cheap |
| Main memory reference | ~100 ns | cheap |
| Compress 1 KB (Snappy) | ~2 µs | cheap |
| Send 1 KB over 1 Gbps network | ~10 µs | noticeable |
| Read 4 KB randomly from SSD | ~150 µs | noticeable |
| Read 1 MB sequentially from memory | ~250 µs | noticeable |
| Round trip within same datacenter | ~0.5 ms | the unit of design |
| Read 1 MB sequentially from SSD | ~1 ms | expensive |
| Disk seek (spinning) | ~10 ms | expensive |
| Read 1 MB sequentially from disk | ~20 ms | expensive |
| Round trip CA → Netherlands → CA | ~150 ms | architectural |
Three consequences you should be able to state instantly: memory is ~100× faster than SSD, SSD is ~20× faster than spinning disk, and a cross-continent round trip costs more than 300 in-datacenter round trips. That last one is why you put a CDN in front of anything a global audience reads, and why chatty microservice calls across regions are an anti-pattern.
Powers of two, and the only conversions you need
The estimation recipe
Always four quantities, in this order. Round aggressively — 86,400 becomes 100,000, and nobody objects.
- QPS = DAU × actions-per-user-per-day ÷ 100,000. Then peak QPS = 2–3 × average (state which multiplier you're using).
- Storage = writes/day × bytes/write × retention days × replication factor.
- Bandwidth = QPS × bytes/request, split into ingress and egress.
- Memory for cache = apply the 80/20 rule — cache the 20% of data serving 80% of reads, for a hot window (often 1 day).
Worked example — Twitter-scale
End every estimation with a conclusion sentence that names the binding constraint. "The writes are trivial; the design problem is read fan-out." That single sentence proves you understood why you did the math, and it sets the agenda for the rest of the interview.
Reference capacities of one machine
| Component | A single commodity node handles roughly |
|---|---|
| Nginx / Envoy proxy | 50k–100k req/s (static), ~10k–20k with TLS termination and routing |
| JVM app server (simple JSON API) | 5k–20k req/s per node, hundreds if it does heavy work per request |
| Redis (single node) | ~100k ops/s, up to 500k with pipelining; 64–256 GB RAM |
| PostgreSQL / MySQL (single node) | ~5k–15k simple reads/s, ~1k–5k writes/s; comfortable to a few TB |
| Cassandra node | ~10k writes/s; scales linearly by adding nodes |
| Kafka broker | ~100 MB/s–1 GB/s per broker, millions of msgs/s per cluster |
| Modern NVMe SSD | ~500 MB/s–3 GB/s sequential, 100k+ random IOPS |
| 1 Gbps NIC | ~125 MB/s; 10 Gbps ~ 1.25 GB/s |
Computing 47,382 QPS to three significant figures. Nobody cares. But failing to notice that 47,000 QPS × 2 KB = 94 MB/s exceeds a 1 Gbps NIC costs you the round. Precision is worthless; noticing which resource saturates first is everything.
Practice — Estimate storage and QPS for a URL shortener with 100M new URLs/day, 10:1 read ratio, 10-year retention.estimation
Write QPS: 100M / 10⁵ ≈ 1,000/s, peak ×3 ≈ 3,000/s.
Read QPS: 10,000/s average, ~30,000/s peak.
Storage per record: short key (7 bytes) + long URL (~100 bytes avg) + userId (8) + createdAt (8) + expiry (8) ≈ 130 bytes; round to 200 bytes with index overhead.
Total records: 100M × 365 × 10 = 365 billion. That is the number that should alarm you.
Total storage: 365B × 200 B ≈ 73 TB, ×3 replication ≈ 220 TB.
Conclusion: 73 TB will not fit on one node, so we must shard — but the access pattern is a pure key lookup with no joins and no range scans, so a hash-partitioned key-value store is the natural fit, and 30k reads/s is trivially absorbed by cache since redirects are extremely skewed (a tiny fraction of links get most traffic).
Bonus point: question the premise. 365 billion links over 10 years is almost certainly wrong for a real product; propose a TTL (e.g. links expire after 2 years unless renewed) and storage drops by 80%. Interviewers love a candidate who pushes back on a requirement with a number.
Practice — A service handles 40,000 req/s. Each response is 50 KB. Can one 10 Gbps datacenter uplink serve it?estimation
40,000 × 50 KB = 2,000,000 KB/s = 2 GB/s = 16 Gbps. A single 10 Gbps link cannot. Options in order of preference: (1) push the payload to a CDN so origin egress collapses — usually the right answer for static or semi-static content; (2) compress (gzip/brotli on JSON often gives 5–10×, bringing you to ~2–3 Gbps); (3) shrink the payload — 50 KB per response usually means you are over-fetching, so add field selection or pagination; (4) only then add uplinks/servers. Notice the ordering: change the data before you buy the hardware.
Networking essentials for designers
What actually happens when a user hits your URL
DNS as a design tool
DNS is not just name resolution — it is your first load balancer and your disaster-recovery switch.
- Round-robin DNS — multiple A records; free, dumb, no health awareness, and clients cache aggressively.
- GeoDNS / latency-based routing — return the IP of the nearest healthy region. This is how you get users to the right datacenter.
- Failover via TTL — a 60-second TTL means a regional outage can be rerouted in about a minute. Long TTLs (24 h) are cheaper but pin you during incidents. State the trade-off in interviews: low TTL = fast failover, more DNS queries; high TTL = fewer queries, slow recovery.
- Anycast — one IP announced from many locations; BGP routes the user to the topologically nearest. Used by CDNs and public resolvers.
TCP vs UDP, and when the answer is genuinely UDP
| TCP | UDP | |
|---|---|---|
| Guarantees | Ordered, reliable, flow + congestion controlled | None — fire and forget |
| Setup | 3-way handshake (1 RTT before data) | No handshake |
| Head-of-line blocking | Yes — one lost packet stalls the stream | No |
| Use it for | Almost everything: HTTP, DB connections, RPC | DNS, real-time voice/video (WebRTC), game state, metrics (StatsD), QUIC's substrate |
QUIC / HTTP-3 is the interesting modern answer: it rebuilds reliability and ordering on top of UDP, per-stream, which removes TCP's head-of-line blocking and folds the transport + TLS handshake into a single round trip. Also it survives IP changes (phone switching Wi-Fi → cellular) because connections are identified by a connection ID, not the 4-tuple. That is a genuinely strong thing to mention in a mobile-heavy design.
HTTP versions — what changed and why you care
| Version | Key property | Design implication |
|---|---|---|
| HTTP/1.1 | One request in flight per connection; browsers open ~6 per host | Domain sharding and sprite-sheets were hacks around this |
| HTTP/2 | Multiplexed streams over one TCP connection, header compression, server push | Those hacks become harmful; TCP head-of-line blocking remains |
| HTTP/3 | QUIC over UDP; per-stream reliability, 1-RTT setup, connection migration | Best for lossy mobile networks; the default at large CDNs today |
Idempotency and HTTP methods
Safe = no side effects (GET, HEAD). Idempotent = doing it twice equals doing it once (GET, PUT, DELETE — not POST). This matters because retries are mandatory in distributed systems, and retrying a non-idempotent POST double-charges a customer. The fix is an idempotency key: the client generates a UUID, sends it as a header, and the server stores (key → result) so a replay returns the original response instead of re-executing. Stripe's API is the canonical example; say so.
CDN — the highest-leverage box on the diagram
A CDN is a globally distributed reverse-proxy cache. It reduces latency (content is physically closer), reduces origin load (cache hit ratio of 90%+ is normal), and absorbs traffic spikes and some DDoS.
- Push CDN — you upload content to the edge proactively. Good for small, rarely-changing catalogues.
- Pull CDN — the edge fetches from origin on first miss and caches per TTL. The default for almost everything.
- Cache key — usually the URL. Beware of caching a personalised response under a shared key; use
Varyheaders or put user-specific content on a separate, uncached path. - Invalidation — either purge by URL/tag (slow, rate-limited) or never invalidate: use content-hashed filenames (
app.9f3c2a.js) with a one-year TTL. The second option is strictly better and is what every serious frontend build does. - Signed URLs — time-limited tokens so private media can still be served from the edge. This is how YouTube/Netflix serve authenticated video without the origin touching the bytes.
Practice — Why is the "seconds to first byte" of a video stream a CDN problem and not a database problem?networking
The database only returns metadata and a manifest — a few hundred bytes, sub-millisecond from cache. The perceived start time is dominated by (a) DNS + TCP + TLS to a server, (b) the distance to whatever serves the first video segment, and (c) the size of that first segment. All three are CDN concerns. The engineering fixes are: edge-terminate TLS close to the user, use HTTP/3 to cut handshake RTTs, and use adaptive bitrate with a small initial segment (2 s at low bitrate) so playback can start before the high-quality data arrives. Mentioning that first-segment size is a deliberate trade-off between startup latency and initial visual quality is a strong depth signal.
How services talk: REST, RPC, GraphQL, and push
Request/response styles
| REST | gRPC | GraphQL | |
|---|---|---|---|
| Transport / format | HTTP/1.1 or 2 + JSON | HTTP/2 + Protobuf (binary) | HTTP + JSON, single endpoint |
| Contract | OpenAPI (optional, often drifts) | .proto — enforced, code-generated | Typed schema, introspectable |
| Strength | Universal, cacheable by HTTP semantics, debuggable with curl | Fast, small payloads, streaming both ways, strict typing | Client asks for exactly the fields it needs; kills over-fetching |
| Weakness | Chatty; over- and under-fetching; versioning pain | Not browser-native (needs grpc-web proxy); binary is harder to debug | Hard to cache at HTTP layer; N+1 resolver problem; complex queries can DoS you |
| Use when | Public APIs, third-party integrations, simple CRUD | Internal service-to-service, high throughput, polyglot backends | Many different clients (iOS/Android/web) with different field needs |
"REST at the edge, gRPC between internal services." This is what most large systems actually do: a public REST/GraphQL gateway for clients, Protobuf-over-HTTP/2 internally where you control both ends and care about latency and payload size.
Server-to-client push — the four options
| Technique | How | Cost | Right for |
|---|---|---|---|
| Short polling | Client asks every N seconds | Wasteful; N/2 average latency | Nothing, mostly. Acceptable for low-frequency status checks |
| Long polling | Server holds the request open until data or timeout | One held connection per client; simple, works everywhere | Fallback path, moderate update rates |
| SSE (Server-Sent Events) | One long-lived HTTP response, server streams text events | Cheap, auto-reconnect built in, one-directional | Live feeds, notifications, progress bars, LLM token streaming |
| WebSocket | HTTP Upgrade → full-duplex TCP | Stateful connection to pin and scale; needs its own LB story | Chat, multiplayer, collaborative editing, trading |
The rule: if the client only needs to receive, use SSE — it is far simpler operationally. If both sides send frequently, use WebSocket. If you only need updates every 30+ seconds, polling is fine and you should say so rather than reflexively reaching for sockets.
The hard part of WebSockets is not the protocol, it's the state
A WebSocket pins a user to one specific server. That breaks the "stateless service" assumption everything else relies on, and creates four problems you must be ready to discuss:
- Routing — to send a message to user B, you must know which gateway node holds B's connection. Solution: a session registry (Redis:
userId → nodeId, with TTL and heartbeat). - Delivery across nodes — node A must hand the message to node B. Solution: a pub/sub bus (Redis Pub/Sub, Kafka, or NATS) with a topic per node or per user.
- Load balancing — connections are long-lived, so a new node gets zero traffic after a scale-out. You need connection draining and sometimes forced rebalancing.
- Deploys — every deploy drops every connection. Clients need reconnect with exponential backoff and jitter, or you get a thundering herd that takes down the freshly deployed cluster.
Practice — WhatsApp has 2 billion users. Do you keep 2 billion WebSockets open?real-time
No — you keep open connections only for currently active clients, which is a small fraction (say 5–10% concurrent). Even so, that is 100–200M concurrent connections. Two things make it feasible: (1) a single tuned Linux box can hold ~500k–1M idle connections (the limits are file descriptors, ephemeral ports per destination, and ~10–50 KB of kernel + app memory per socket), so you need thousands of gateway nodes, not millions; (2) for backgrounded mobile apps you do not hold a socket at all — you fall back to push notifications (APNs/FCM), which is one connection from the device to Apple/Google that the OS already maintains for every app. The design answer is therefore a hybrid: WebSocket for foreground, platform push for background, and a durable message store so nothing is lost while offline.
Scaling: statelessness and load balancers
Vertical vs horizontal
Vertical (bigger machine) is underrated: it is free of distributed-systems complexity, and modern cloud instances go to hundreds of cores and terabytes of RAM. It fails on two counts — there is a ceiling, and one machine is a single point of failure. Horizontal (more machines) has no ceiling and gives you redundancy, but every problem you had becomes a distributed problem: coordination, consistency, partial failure.
"I'd scale this vertically first, because a single Postgres node with 64 cores handles our projected 5k writes/sec comfortably, and it saves us the entire class of distributed-consistency bugs. I'd design the schema so that sharding by tenantId later is possible, but I wouldn't shard on day one." Senior engineers defer complexity; junior engineers accumulate it.
Statelessness is the enabling property
A stateless service stores nothing about a client between requests, so any node can serve any request, and you can add or kill nodes freely. State does not disappear — you relocate it: session data to Redis or a signed cookie/JWT, uploaded files to blob storage, in-flight jobs to a queue.
Load balancers: L4 vs L7
| L4 (transport) | L7 (application) | |
|---|---|---|
| Sees | IP + port only | Full HTTP: path, headers, cookies, method |
| Can do | Forward TCP/UDP extremely fast; millions of conns | Path routing, TLS termination, header rewriting, retries, rate limiting, canary by header |
| Cost | Lowest latency, minimal CPU | Must decrypt and parse; more CPU |
| Examples | AWS NLB, IPVS, HAProxy in TCP mode | AWS ALB, Nginx, Envoy, Traefik |
Real architectures use both: an L4 layer (or Anycast) absorbs raw connections and DDoS, then hands to an L7 layer that does routing and TLS.
Balancing algorithms — and the one that is usually right
- Round robin — even distribution, ignores that request cost varies wildly.
- Weighted round robin — for heterogeneous hardware.
- Least connections — good proxy for load when request durations vary.
- Least response time — best general-purpose choice; adapts to a degraded node automatically.
- IP hash / consistent hash — sends the same client (or key) to the same node; used for cache locality, not for fairness.
- Power of two choices — pick two nodes at random, send to the less loaded one. Nearly as good as global least-connections but needs no global state, so it scales to huge fleets. Mentioning this is a genuine depth signal.
Two kinds: a liveness check (is the process up?) and a readiness check (can it serve traffic — DB reachable, caches warm?). Failing to distinguish them causes the classic incident where a node restarts, is marked healthy immediately, receives full traffic with a cold cache and empty connection pool, times out, gets marked unhealthy, and flaps. Fixes: slow-start / ramp-up on new nodes, and readiness checks that actually exercise a dependency.
Scaling the whole picture, in the order you should propose it
Practice — Your read replicas are added, but users complain that after posting, their own post is missing. Why, and three fixes?replication
Why: replication lag. The write went to the leader; the subsequent read was served by a follower that had not yet applied it. This is the read-your-own-writes violation.
Fix 1 — read from the leader for a window. After a user writes, route that user's reads to the leader for N seconds (track a timestamp in their session). Simple, effective, slightly increases leader load.
Fix 2 — monotonic reads via pinning. Hash the user to a specific replica so they never go backwards in time, even if that replica is behind.
Fix 3 — write-through the client. Return the created object in the write response and have the client render it optimistically; or pass the write's log-sequence-number back and have the read wait until the replica has caught up to that LSN (this is how "causal consistency tokens" work in DynamoDB/Cosmos-style systems).
Naming the failure as "read-your-own-writes consistency" by its proper term is worth as much as the fix.
Databases I: models, indexes, storage engines
"SQL or NoSQL?" is the wrong question
The right question is what is my access pattern, and what does the storage engine make cheap? Answer these five and the database chooses itself:
- Do I query by a single key, or do I need flexible ad-hoc queries and joins?
- Is the write volume beyond what one node can take?
- Do I need multi-object atomic transactions?
- Is the schema stable, or genuinely heterogeneous per record?
- Do I need range scans / sorting, or only point lookups?
| Family | Examples | Data model | Choose when |
|---|---|---|---|
| Relational | PostgreSQL, MySQL | Tables, rows, foreign keys | Default. Ad-hoc queries, joins, ACID transactions, moderate scale. Postgres to a few TB and ~10k writes/s is boring and correct. |
| Wide-column | Cassandra, ScyllaDB, HBase, Bigtable | Partition key + clustering columns | Huge write volume, known access pattern, time-series or per-user timelines. No joins. Query shape is fixed at table-design time. |
| Key-value | DynamoDB, Redis, Riak | Opaque blob by key | Sessions, caches, feature flags, anything strictly point-lookup. |
| Document | MongoDB, Couchbase | JSON documents | Records that vary in shape and are read whole (product catalogues, CMS). |
| Graph | Neo4j, JanusGraph | Nodes + edges | Multi-hop traversals: friend-of-friend, fraud rings, dependency graphs. |
| Search | Elasticsearch, OpenSearch | Inverted index | Full-text, faceting, relevance ranking. A secondary index, never a system of record. |
| Columnar / OLAP | ClickHouse, Snowflake, Redshift, Druid | Column-oriented | Analytics over billions of rows: aggregations across few columns. |
| Time-series | InfluxDB, TimescaleDB, Prometheus | (metric, tags, ts) → value | Metrics and IoT: append-only, time-ordered, aggressive downsampling. |
NoSQL stores usually buy scale by giving up two things: joins and cross-partition transactions. In exchange you get linear horizontal scaling. That means you denormalise and often store the same data multiple times, once per query pattern. If you propose Cassandra, you must be ready to say "and I'd store the timeline twice, once keyed by userId and once by videoId" — otherwise the interviewer knows you have only read the marketing page.
Indexes: the single highest-value DB topic
An index is a separate data structure that maps a column value to row locations, so a lookup becomes O(log n) instead of O(n). The cost is that every write must update every index — indexes make reads faster and writes slower, and consume disk.
B+ tree (Postgres, MySQL/InnoDB — the read-optimised default)
LSM tree (Cassandra, RocksDB, LevelDB — the write-optimised default)
"B-trees give you read-optimised, in-place updates; LSM trees give you write-optimised, append-only storage with read amplification and compaction cost. Since our workload is 50k writes/sec of immutable events, I'd take the LSM engine." Naming the storage engine rather than the product is the difference between a mid and a senior answer.
Index types you should be able to name
- Clustered index — the table itself is stored in index order (InnoDB primary key). One per table; range scans on it are extremely fast. A random UUID primary key destroys this by causing page splits everywhere — use an ordered ID (see H3, Snowflake) or UUIDv7.
- Secondary (non-clustered) index — maps value → primary key, so a lookup costs two traversals unless the index is covering.
- Composite index —
(a, b, c)serves queries ona,(a,b),(a,b,c)— the leftmost-prefix rule — but not onbalone. This is one of the most-asked DB interview questions. - Covering index — includes all columns the query needs, so the engine never touches the heap. The single best query optimisation available.
- Hash index — O(1) equality only, no ranges.
- Bitmap index — low-cardinality columns in analytics stores.
- Inverted index — token → list of documents. The basis of search (F14).
- Geospatial — R-tree, quadtree, geohash, S2 cells (see H12).
Practice — A query is slow. Walk me through your diagnosis.databases
1. EXPLAIN ANALYZE it. Look for sequential scans on large tables, nested loops over big row counts, and — most importantly — a large gap between estimated and actual rows, which means stale statistics.
2. Check for a missing or unusable index. Common causes of an unusable index: a function on the column (WHERE lower(email) = ? needs a functional index), an implicit type cast, a leading wildcard LIKE '%foo', or an OR across columns.
3. Check row volume. If the query legitimately touches millions of rows, no index saves you — you need pre-aggregation, a materialised view, or a different store (columnar).
4. Check locking/contention — a slow query may be waiting, not working. Look at lock waits and the isolation level.
5. Check the plan cache and parameter sniffing — the same SQL can be fast for one parameter and catastrophic for another.
6. Only then consider denormalisation or caching. Caching a slow query hides the problem and makes invalidation your new problem.
Databases II: transactions, ACID, isolation
ACID, precisely
- Atomicity — all operations in the transaction happen, or none do. Implemented with a write-ahead log and rollback.
- Consistency — the database moves from one valid state to another, respecting constraints. This is largely the application's responsibility; it is the weakest letter and the one people over-explain.
- Isolation — concurrent transactions do not corrupt each other. The interesting letter.
- Durability — once committed, it survives a crash. Implemented by fsync-ing the WAL, and in distributed systems by replicating to a quorum before acknowledging.
Isolation levels and the anomalies they prevent
| Level | Dirty read | Non-repeatable read | Phantom read | Cost |
|---|---|---|---|---|
| Read Uncommitted | possible | possible | possible | lowest |
| Read Committed (Postgres default) | prevented | possible | possible | low |
| Repeatable Read (MySQL default) | prevented | prevented | possible* | medium |
| Serializable | prevented | prevented | prevented | highest |
*InnoDB's Repeatable Read prevents most phantoms using next-key (gap) locks; Postgres's Repeatable Read is snapshot isolation and prevents phantoms for reads but allows write skew.
Booking systems (H13) and any "check-then-act" logic hit write skew. The fixes: (a) SERIALIZABLE isolation, (b) explicit locking with SELECT ... FOR UPDATE on the rows you checked, (c) materialise the conflict — introduce a row that both transactions must write to (e.g. a seat_hold row with a unique constraint on (showId, seatId)), so the DB turns the logical conflict into a physical one it can detect.
Pessimistic vs optimistic concurrency
Note the escape hatch: after N attempts you must fail loudly, not loop forever. Unbounded retries under contention are how a hot SKU takes down a checkout service.
Deadlocks
Four Coffman conditions must all hold: mutual exclusion, hold-and-wait, no preemption, circular wait. Databases usually detect cycles and kill a victim transaction, returning a retryable error. The practical prevention is lock ordering: always acquire locks on rows in a deterministic order (e.g. ascending primary key), so a cycle cannot form. In a money-transfer design, say: "I lock the two accounts in ascending account-id order" — it is a one-line answer that shows you have debugged this in real life.
Practice — Two users transfer money to each other simultaneously. Show the deadlock and fix it.transactions
Deadlock: T1 (A→B) locks A then waits for B. T2 (B→A) locks B then waits for A. Circular wait; both block until the DB kills one.
Fix 1 — ordered locking: both transactions lock min(accountId) first, then max(accountId). No cycle is possible, so no deadlock ever occurs.
Fix 2 — single-statement atomicity where possible: UPDATE accounts SET balance = balance - 100 WHERE id = ? AND balance >= 100 — the check and the change happen in one statement, which the DB executes atomically, and rowsAffected = 0 means insufficient funds.
Fix 3 — event sourcing / ledger: never update balances in place. Append immutable double-entry ledger rows and compute the balance as a sum (with periodic snapshots). This eliminates the contention entirely, is auditable, and is what real payment systems do (see H14).
Replication and partitioning
These are the two orthogonal axes of distributing data, and mixing them up is a common mistake. Replication = the same data on multiple nodes (for availability and read throughput). Partitioning/sharding = different data on different nodes (for write throughput and total capacity). Real systems do both: shard the data, replicate each shard.
Replication topologies
| Topology | How | Trade-off |
|---|---|---|
| Single-leader | All writes to leader; followers replay its log | Simple, no write conflicts. Leader is a write bottleneck and a failover event. The default; choose this unless you have a reason. |
| Multi-leader | Multiple nodes accept writes and replicate to each other | Great for multi-region write latency and offline clients. You must solve write conflicts (LWW, CRDT, app-level merge). |
| Leaderless (Dynamo-style) | Client writes to W nodes, reads from R nodes | Highly available, no failover. Needs quorums, read-repair, anti-entropy, and version reconciliation. |
Synchronous vs asynchronous replication
Sync: the leader waits for the follower to acknowledge before confirming the write. No data loss on leader failure, but a slow follower slows every write, and if it dies, writes stall. Async: the leader confirms immediately. Fast and available, but a leader crash loses the un-replicated tail. Semi-sync (wait for exactly one follower) is the practical middle ground and what most production MySQL/Postgres setups use.
"Async replication with typical lag of tens of milliseconds is fine for feed reads, but for the balance shown on a payment confirmation screen I'd read from the leader, because a stale balance is a support ticket." Tying the consistency choice to a concrete product consequence is exactly what senior means here.
Failover, and why it is dangerous
- Detect — usually a timeout on heartbeats. Too short → spurious failovers; too long → extended downtime.
- Elect — a consensus algorithm (Raft) or an external coordinator (ZooKeeper/etcd) picks the most up-to-date follower.
- Reconfigure — clients and remaining followers must learn the new leader.
- Split-brain — the old leader may still think it is the leader and accept writes. Prevented by fencing tokens: every leadership term gets a monotonically increasing number, and storage rejects writes carrying an old token.
Partitioning strategies
| Strategy | Mechanism | Good | Bad |
|---|---|---|---|
| Range | Keys A–F → shard 1, G–M → shard 2… | Range scans stay on one shard; ordered | Hotspots — sequential keys (timestamps!) all hit the newest shard |
| Hash | hash(key) % N | Even distribution | Range scans must fan out to all shards; adding a node reshuffles everything |
| Consistent hash | Ring with virtual nodes (F9) | Adding/removing a node moves only 1/N of keys | Slightly more complex; still no range locality |
| Directory / lookup | A metadata service maps key → shard | Total flexibility; rebalance any key | The lookup service is a new SPOF and an extra hop; must be cached |
| Geographic | Partition by region | Data residency (GDPR), low local latency | Cross-region queries are painful; uneven region sizes |
Choosing a shard key — the highest-stakes decision in the design
A good shard key has three properties: high cardinality (many distinct values), even distribution (no value dominates), and query alignment (your most common query can be answered from a single shard). You will usually be asked to defend it.
The classic: Justin Bieber's userId, or a viral tweet, makes one shard 100× hotter than the rest. Three fixes, in increasing order of ugliness: (1) cache the hot key — usually enough, because hot means read-hot; (2) salt the key — write to key#0…key#9 and read all ten, spreading the write load 10× at the cost of a 10-way scatter-gather; (3) special-case it — Twitter genuinely does not fan out celebrity tweets, merging them at read time instead (see H7).
Resharding without downtime
- Double-write to old and new topology while backfilling historical data.
- Verify by comparing reads from both (shadow reads, diff the results, alert on mismatch).
- Flip reads to the new topology behind a feature flag, per-tenant if possible.
- Stop double-writing; decommission the old shards after a safety window.
Alternatively, pre-shard: create 1,024 logical partitions on day one and map many of them onto few physical nodes. Growing then means moving whole logical partitions — no key remapping at all. This is what Vitess, Citus, and most mature systems do, and proposing it unprompted is a strong signal.
Practice — You sharded by userId. Product now wants "search all posts by hashtag." What breaks and what do you do?sharding
What breaks: hashtag search has no userId, so every query must scatter to all N shards and gather — latency becomes the slowest shard's latency (tail amplification), and cost scales with N per query.
Answer: do not force the primary store to serve a query pattern it is not partitioned for. Build a secondary index optimised for that access path: either (a) a separate table/store partitioned by hashtag (hashtag → postIds, written asynchronously from a change stream), or (b) push the data into a search cluster (Elasticsearch) which is designed for inverted-index lookups and relevance ranking.
Then name the consequence honestly: the secondary index is eventually consistent, typically lagging by hundreds of milliseconds to seconds, and you need a reconciliation job to repair drift. The general principle — one store per access pattern, joined by a change stream — is the answer to a whole family of interview questions.
Consistent hashing
The problem it solves
The ring
"Adding or removing one node relocates about K/N keys instead of all of them. Virtual nodes fix the uneven distribution you get with a small number of physical servers, and they also spread a departed node's load across the whole cluster rather than onto one neighbour. The replica set for a key is just the next N distinct physical nodes clockwise — which is exactly how Dynamo and Cassandra place replicas."
Where it shows up
- Distributed caches (Memcached clients, Redis Cluster uses 16,384 fixed hash slots — a pre-sharding variant of the same idea).
- Dynamo-style databases: Cassandra, DynamoDB, Riak, ScyllaDB.
- Sharded WebSocket gateways and rate limiters (route a user consistently to the node holding their counter).
- Load balancers doing cache-affinity routing (
consistent hashingupstream mode in Nginx/Envoy).
Practice — What's the weakness of consistent hashing, and what's the modern alternative?deep-dive
Weaknesses: (1) even with virtual nodes the distribution is only statistically even — real load imbalance of 10–20% is common because key popularity is not uniform even when key placement is; (2) it gives you no control over placement, so you cannot enforce constraints like "no two replicas in the same rack"; (3) memory for the ring grows with V × N.
Alternatives worth naming: Rendezvous (HRW) hashing — compute hash(key, node) for every node and take the max; same minimal-disruption property, no ring, and trivially supports weights, at O(N) per lookup. Jump consistent hash (Google) — O(ln N) time, no memory, but only supports adding/removing at the tail. Bounded-load consistent hashing — the ring plus a cap on any node's share, spilling overflow to the next node; this is what Vimeo/Google use to fix the popularity-skew problem. Mentioning bounded-load or rendezvous marks you as someone who has read past the standard blog posts.
CAP, PACELC, quorums, and consistency models
CAP, stated correctly
The common phrasing "pick two of three" is wrong and interviewers notice. Network partitions are not a choice — on a real network they will happen. So the theorem really says: when a partition occurs, you must choose between consistency and availability. When there is no partition, you can have both.
- CP — during a partition, refuse to serve requests that cannot be made consistent. Examples: ZooKeeper, etcd, HBase, Spanner, MongoDB (with majority writes). Choose for balances, inventory, locks, leader election.
- AP — during a partition, keep serving, accept divergence, reconcile later. Examples: Cassandra, DynamoDB, Riak, DNS. Choose for feeds, likes, product catalogues, session stores, metrics.
If Partition, then Availability or Consistency; Else, Latency or Consistency. The second half is the part you deal with daily: even with a perfectly healthy network, waiting for a quorum across regions costs tens of milliseconds. Every "eventually consistent" design in the real world is chosen for latency, not for partition tolerance. Saying this is one of the fastest ways to sound senior.
Consistency models, strongest to weakest
| Model | Guarantee | Cost / example |
|---|---|---|
| Linearizable (strong) | Every read sees the most recent completed write; the system behaves as if there is one copy | Needs consensus or leader reads. etcd, Spanner, single-node RDBMS. Expensive across regions. |
| Sequential | All nodes see operations in the same order, not necessarily real-time order | Cheaper than linearizable, rarely the explicit target |
| Causal | Operations that are causally related are seen in order by everyone; concurrent ones may differ | The sweet spot for social apps — a reply never appears before its parent. Implemented with version vectors. |
| Read-your-writes | A user always sees their own writes | Session guarantee. Achieved by leader-reads-after-write or sticky routing. |
| Monotonic reads | Time never moves backwards for a given user | Pin a user to a replica. |
| Eventual | If writes stop, replicas converge — eventually | Cheapest and most available. Fine for like counts, view counts, feeds. |
Note that the middle four are session guarantees: they are per-user promises, not global ones, which makes them cheap. Most products need "eventual consistency plus read-your-writes," and being able to say precisely that — rather than "strong" or "eventual" — is the mark of someone who has shipped.
Quorums
When a read returns conflicting versions, the system must reconcile: last-write-wins (simple, silently loses data, needs synchronised clocks), version vectors (detect concurrent writes and hand both to the application — Riak's siblings), or CRDTs (data types whose merge is mathematically guaranteed to converge: G-counters, PN-counters, OR-sets, LWW-registers). Two repair mechanisms keep replicas honest in the background: read repair (fix stale replicas when a read notices divergence) and anti-entropy (periodic Merkle-tree comparison between replicas).
It depends on wall-clock time across machines. Clock skew of even 50 ms silently discards a write that arrived "later" in real time. This is not theoretical — it is a known Cassandra footgun. Fixes: use logical clocks (Lamport timestamps / version vectors), or use a system with tightly bounded clock uncertainty (Spanner's TrueTime, which waits out the uncertainty interval before committing).
Practice — Is Kafka CP or AP? Justify.deep-dive
Kafka is tunable, and the honest answer names the knobs. With acks=all, min.insync.replicas=2 on a replication factor of 3, and unclean.leader.election.enable=false, a partition that loses its in-sync replicas becomes unavailable for writes rather than accepting writes that could be lost — that is CP behaviour. Flip acks=1 and allow unclean leader election, and Kafka will keep accepting writes and may silently lose the tail of the log — AP behaviour. The interview point is that "CP or AP" is usually a configuration decision at the boundary of a subsystem, not an inherent property of a product, and you should state which configuration you are assuming.
Practice — Design a like-counter for a post that gets 500k likes/minute. Consistency choice?applied
Eventual consistency, deliberately. Nobody is harmed by a like count that is 3 seconds stale, and forcing linearizability on a single hot row creates a contention hotspot that will fail at a fraction of this rate.
Design: (1) client sends the like; (2) service appends an event to Kafka partitioned by postId (also gives you idempotency via a (userId, postId) dedupe key so double-taps count once); (3) a stream processor aggregates per-post counts in memory and flushes deltas every second; (4) the count is stored in Redis as an atomic INCRBY and periodically checkpointed to durable storage; (5) reads come from Redis with a short TTL, and the client applies an optimistic +1 locally so the tap feels instant.
Sharding the counter: if a single post is still too hot, split the counter into 100 sub-counters (likes:{postId}:{0..99}), increment a random one, and sum on read — a bounded scatter that turns one hot key into a hundred warm ones. This is the standard "counter sharding" answer and is worth naming explicitly.
Caching
Where caches live
The four cache patterns
| Pattern | Read | Write | Trade-off |
|---|---|---|---|
| Cache-aside (lazy loading) | App checks cache; on miss reads DB and populates cache | App writes DB, then deletes the cache entry | Simplest and most common. Only requested data is cached. First request after a miss is slow; risk of stale entries on race conditions. |
| Read-through | Cache itself fetches from DB on miss | — | Cleaner app code, but you need a cache library/provider that supports loaders. |
| Write-through | Always from cache | Write to cache and DB synchronously | Cache is never stale. Every write pays both latencies, and you cache data nobody may read. |
| Write-back (write-behind) | From cache | Write to cache; flush to DB asynchronously in batches | Fastest writes, absorbs spikes. Data loss if the cache dies before flush — only acceptable for tolerable data (view counts, metrics). |
On a write, invalidate the cache entry rather than writing the new value into it. Updating creates a race: two concurrent writers can interleave so the cache ends up holding the older value permanently. Deleting means the next read re-populates from the source of truth. If you must reduce the miss, use DELETE then optionally re-populate after the DB commit, and accept that a rare miss is far cheaper than permanent staleness.
Eviction policies
- LRU — evict least-recently-used. Good default; implemented as a hash map + doubly linked list (see L15 for full Java).
- LFU — evict least-frequently-used. Better when popularity is stable, but suffers from "cache pollution" by items that were hot long ago; real implementations use windowed/decayed counts.
- FIFO — simple, ignores access patterns, rarely right.
- TTL — not really eviction but the workhorse: bounded staleness with no invalidation logic at all. Always add jitter to TTLs so a million keys written together do not expire together.
- W-TinyLFU (Caffeine's default) — an admission filter using a frequency sketch in front of an LRU-ish window. Near-optimal hit rates; naming it is a strong signal.
The four cache failure modes you must be able to name
| Failure | What happens | Fix |
|---|---|---|
| Thundering herd / stampede | A popular key expires; 10,000 concurrent requests all miss and all hit the DB | Request coalescing (single-flight): the first miss takes a lock and fetches, others wait on the same future. Or probabilistic early recomputation — refresh a key slightly before expiry with probability rising as TTL nears zero. |
| Cache penetration | Requests for keys that do not exist always miss and always hit the DB — trivially weaponised by an attacker | Cache the negative result with a short TTL, and/or front the cache with a Bloom filter of existing keys. |
| Cache avalanche | Many keys expire simultaneously (mass warm-up, same TTL), or the cache cluster restarts cold | TTL jitter; staged warm-up; keep the cache in a separate failure domain; rate-limit or shed load to the DB during recovery. |
| Hot key | One key gets so much traffic it saturates a single Redis shard's CPU/NIC | Replicate the key to a local in-process cache with a 1-second TTL; or split into key:0..N replicas and pick randomly on read. |
Redis, concretely
Redis is single-threaded for command execution (which is why every command is effectively atomic), in-memory, with optional persistence via RDB snapshots and/or an AOF append-only log. Its value is that it ships useful data structures, not just strings:
| Structure | Used for |
|---|---|
String + INCR | Counters, rate limiters, distributed locks (SET key val NX PX 30000) |
Hash | Objects with fields; cheaper than N keys |
Sorted Set (ZSET) | Leaderboards, sliding-window rate limiters, priority queues, time-ordered feeds. The most useful one in interviews. |
List | Simple queues, recent-items lists (LPUSH + LTRIM = capped timeline) |
Set | Unique visitors, tags, membership checks |
HyperLogLog | Approximate cardinality (unique users) in 12 KB with ~0.8% error |
Streams | Kafka-lite: consumer groups, acknowledgements, replay |
Geo | Geohash-backed radius queries (used in H12) |
| Lua scripts | Multi-step atomic operations — the correct way to build a rate limiter |
Practice — Where would you NOT put a cache, and why?judgement
Do not cache when: (1) the data must be strictly correct at read time — account balances during a transfer, remaining inventory at checkout, permission checks for a security decision; (2) the hit rate would be low — caching a long-tail keyspace (e.g. per-user analytics queries that are never repeated) adds memory cost, an extra hop, and an invalidation problem for nothing; (3) the underlying query is already sub-millisecond on an indexed primary key — you are adding a network hop to save nothing; (4) write volume is comparable to read volume, so entries are invalidated before they are ever read.
The general framing: a cache is a correctness liability you accept in exchange for a latency and load benefit. If you cannot state the benefit in numbers, do not take the liability. Saying this out loud distinguishes you from candidates who add Redis to every diagram reflexively.
Queues, logs, and stream processing
Why asynchrony is an architectural tool, not a detail
Putting a queue between a producer and a consumer buys four things: decoupling (they can be deployed and scaled independently), buffering (a traffic spike becomes a longer queue instead of an outage), retry (failed work stays in the queue), and fan-out (one event, many independent consumers). The cost is eventual consistency and a whole new class of bugs: duplicates, reordering, and poison messages.
"The user-facing request should only do what is needed to return a response. Upload the file, write one row, publish an event, return 202 Accepted. Transcoding, thumbnailing, notifying followers, and updating search indexes all happen off the request path." Naming p99 latency as the reason — "so a slow downstream dependency can't hold the user's connection open" — is what makes it an engineering argument rather than a pattern recital.
Queue vs log — the distinction people get wrong
| Message queue (RabbitMQ, SQS) | Distributed log (Kafka, Pulsar, Kinesis) | |
|---|---|---|
| Message lifetime | Deleted after acknowledgement | Retained for a configured period; consumers track an offset |
| Consumers | Compete for messages; work is split | Independent consumer groups each read the whole stream |
| Ordering | Usually best-effort; hard to guarantee with multiple consumers | Strict order within a partition |
| Replay | No | Yes — rewind the offset. Enormously useful for backfills and bug recovery. |
| Routing | Rich (exchanges, topics, headers, priorities, delays) | Simple — partitioned topics |
| Best for | Task distribution: send email, resize image, per-message ack semantics | Event streaming, multi-consumer fan-out, event sourcing, analytics pipelines |
Kafka's mental model in five facts
- A topic is split into partitions; each partition is an append-only, ordered log on disk.
- The partition key decides ordering: all messages with the same key go to the same partition and are therefore strictly ordered relative to each other. Ordering is per-key, never global.
- Parallelism is capped by partition count: within a consumer group, one partition is consumed by at most one consumer. 12 partitions → at most 12 useful consumers.
- Each partition has a leader and followers (ISR = in-sync replicas).
acks=all+min.insync.replicas=2is the durable configuration. - Consumers store an offset. Committing it before processing gives at-most-once; committing after gives at-least-once. There is no free exactly-once.
Delivery semantics and the exactly-once myth
The transactional outbox — how to write to a DB and a queue atomically
You cannot commit a database transaction and publish to Kafka atomically; distributed transactions across them are unavailable in practice. If you write the DB then publish and crash in between, the event is lost. If you publish then write and the write fails, you have published a lie.
Operational realities to mention
- Dead-letter queue (DLQ) — after N failed attempts, move the message aside so one poison message does not block the partition, and alert on DLQ depth.
- Backpressure — if consumers cannot keep up, lag grows without bound. Monitor consumer lag as a first-class SLI; scale consumers up to the partition count, then add partitions.
- Ordering vs parallelism — these are in direct tension. If you need per-user ordering, partition by userId and accept that one slow user's messages block only that partition.
- Delayed / scheduled messages — SQS delay queues, RabbitMQ TTL+DLX, or a Redis ZSET keyed by execution time (see H16).
Practice — Consumer lag is growing at 50k messages/minute. Walk me through your response.operations
Immediate triage: is the producer rate up (a traffic spike or a buggy retry loop), or is consumer throughput down (a slow downstream dependency, a GC problem, a poison message causing repeated retries)? The metrics that answer this are produce rate, consume rate, and per-message processing time.
If consumers are slow: scale out consumers — but only up to the partition count, since extra consumers in a group sit idle. If already at the cap, add partitions (note: this changes key→partition mapping and therefore breaks ordering for in-flight keys, so it is not free). Also check whether processing can be batched, or whether a synchronous downstream call can be made async.
If a dependency is down: apply a circuit breaker so the consumer fails fast rather than blocking on timeouts, and route to a retry topic with backoff instead of hammering.
Load-shedding option: for non-critical streams, it is legitimate to drop or sample messages to protect the system — but that must be an explicit product decision, not an accident.
Prevention: alert on lag trend, not absolute value; keep consumers at ~50% utilisation so there is headroom to catch up; ensure replay from an earlier offset is a rehearsed operation.
Storage: blobs, files, and analytics stores
Block vs file vs object
| Block (EBS, SAN) | File (NFS, EFS) | Object (S3, GCS) | |
|---|---|---|---|
| Abstraction | Raw disk you format | Hierarchical directories | Flat key → blob + metadata |
| Access | Attached to one VM | Shared mount, POSIX | HTTP API |
| Scale | TBs | TBs–PBs | Effectively unlimited |
| Use for | Database data files, boot volumes | Shared config, legacy apps, build artifacts | Images, video, backups, data lake, logs |
Never store blobs in your relational database. Store the bytes in object storage and keep only the key + metadata in the DB. Then use pre-signed URLs so the client uploads directly to the object store and downloads directly from the CDN — your application servers never touch the bytes. This single decision removes bandwidth, memory pressure, and timeout problems from your API tier, and stating it early in any media design (H10, H11, H14) marks you as having built one.
How object storage achieves durability
Naive replication (3 copies) costs 3× storage. Erasure coding splits an object into k data fragments plus m parity fragments, spread across failure domains; any k of the k+m fragments reconstruct the object. A common (6,3) scheme survives 3 simultaneous failures at 1.5× storage instead of 4×. The trade-off: reconstruction on read is CPU-heavy and slower, so hot data often uses replication while cold data uses erasure coding. Storage classes (Standard → Infrequent Access → Glacier) formalise exactly this trade with lifecycle policies.
OLTP vs OLAP — and why you need both
| OLTP (row store) | OLAP (column store) | |
|---|---|---|
| Query shape | "Give me order 12345" — few rows, many columns | "Sum revenue by region for 2 years" — billions of rows, 3 columns |
| Layout | All of a row's columns stored together | Each column stored contiguously → huge compression, and only the needed columns are read from disk |
| Writes | Frequent small updates | Bulk append; updates are expensive or unsupported |
| Examples | Postgres, MySQL, DynamoDB | ClickHouse, BigQuery, Snowflake, Druid, Redshift |
The link between them is a pipeline: CDC or events from OLTP → Kafka → object storage (Parquet on a data lake) → warehouse. When an interviewer asks "how would the analytics team query this?", the answer is never "they'll query the production database."
Lambda vs Kappa architecture
Lambda runs a batch layer (accurate, slow, reprocesses everything) alongside a speed layer (approximate, real-time), and merges them at query time. Its flaw is maintaining the same logic twice. Kappa drops the batch layer: everything is a stream, and "reprocessing" means replaying the log from an earlier offset through a new version of the job. Kappa is simpler and possible because durable logs made replay cheap. Mention both in analytics designs (H15) and say why you'd pick Kappa.
Search and the inverted index
Why LIKE '%query%' is not search
A leading wildcard cannot use a B-tree index, so it is a full table scan; it has no notion of relevance, no stemming, no typo tolerance, and no way to rank. Search needs a different index structure.
Ranking: TF-IDF and BM25
Term frequency — a term appearing often in a document suggests relevance. Inverse document frequency — a term appearing in every document (like "the") carries no information, so it is down-weighted. BM25 refines this with saturation (the 20th occurrence of a word adds much less than the 2nd) and length normalisation (a match in a short title beats one in a long body). You do not need the formula; you need to say why those two corrections exist. Modern systems then add a second stage: retrieve ~1000 candidates with BM25 or a vector index, then re-rank with a learned model using click-through, recency, and personalisation features.
Elasticsearch architecture, briefly
- An index is split into shards (each a self-contained Lucene index); each shard has replicas. Shard count is fixed at creation — a classic operational trap.
- A query scatters to all shards, each returns its local top-k, and the coordinator gathers and merges. Latency is bounded by the slowest shard.
- Lucene segments are immutable: writes go to a buffer, are flushed into new segments, and merged in the background. Deletes are tombstones. This means near-real-time, not real-time — the default refresh interval is 1 second.
- It is not a system of record. Keep the source of truth in your primary DB and stream changes into the index (via the outbox/CDC pattern from F12).
"Search is a derived index, so it is eventually consistent with the database. If a user creates a listing and immediately searches for it, they may not find it for a second. If the product cannot accept that, I'd read-through to the primary store for the user's own items and merge with search results." Volunteering the staleness and its mitigation is worth more than the architecture itself.
Reliability: timeouts, retries, breakers, shedding
Availability arithmetic
| Availability | Downtime per year | Per month |
|---|---|---|
| 99% ("two nines") | 3.65 days | 7.2 hours |
| 99.9% | 8.76 hours | 43 minutes |
| 99.99% | 52.6 minutes | 4.3 minutes |
| 99.999% | 5.26 minutes | 26 seconds |
Two facts that follow. Serial dependencies multiply: a request touching five services each at 99.9% has 99.5% availability — worse than any component. Redundancy adds nines: two independent components at 99% give 99.99% if failures are truly independent, which they usually are not (shared power, shared network, shared deploy, shared bug). Say that caveat out loud.
The patterns, in the order you should apply them
- Timeouts on everything. A missing timeout is the single most common cause of cascading failure: threads pile up waiting on a dead dependency until the pool is exhausted and a healthy service goes down. Set the timeout from the p99 of the dependency, not from a round number.
- Retries with exponential backoff and jitter — and only for idempotent, retryable errors (timeouts, 503, connection reset). Never retry a 400. Cap total attempts, and budget retries so a deep call chain does not multiply: 3 retries at each of 4 layers is 81 requests for one user action. Use a retry budget (e.g. retries may not exceed 10% of traffic).
- Circuit breaker — after a failure threshold, stop calling the dependency entirely for a cooldown, fail fast, and periodically let one probe through (half-open). This converts a slow cascading failure into a fast, contained one.
- Bulkhead — separate thread pools/connection pools per dependency so a slow one cannot consume all capacity. Named after ship compartments.
- Load shedding — when overloaded, reject the cheapest-to-reject requests early (return 429/503 at the edge) rather than degrading everything. Prioritise: shed anonymous browsing before shedding checkout.
- Graceful degradation — serve stale cache, hide the recommendations carousel, disable non-essential features. Decide the degradation ladder in advance and say what it is.
- Hedged requests — for latency-sensitive reads, send a duplicate request to a second replica if the first has not answered by p95, and take whichever returns first. Cuts tail latency dramatically at a few percent extra load.
The most common self-inflicted outage: a dependency slows down, every client retries, load triples, the dependency dies completely, and it cannot recover because the retry traffic never stops even after it restarts. Defences: jittered backoff, retry budgets, circuit breakers, and — critically — the server telling clients to back off (429 with Retry-After). Mentioning that the server should participate in backpressure, not just the client, is a strong depth signal.
Practice — Service A calls B calls C. C's p99 goes from 20 ms to 2 s. Trace the failure and prevent it.reliability
The cascade: B's threads block on C. B's thread pool fills. B stops responding to A within its timeout — or worse, B has no timeout and holds connections forever. A's threads then block on B. A's pool fills. The user-facing service is down, and monitoring shows A as the failure even though C is the cause.
Prevention, layer by layer: (1) every call has a timeout tighter than the caller's own budget — if A must answer in 500 ms, B's budget is 300 ms and C's is 150 ms, and each hop passes the remaining deadline downstream (deadline propagation); (2) a bulkhead so calls to C use a bounded, separate pool, capping the blast radius; (3) a circuit breaker on C so B fails fast and returns a degraded response; (4) A serves stale cached data or hides the feature C powers; (5) async where possible — if C's result is not needed for the response, publish an event instead of calling it.
The framing that wins: "The goal isn't to prevent C from being slow — it's to make sure C being slow doesn't turn into A being down."
Coordination: consensus, locks, sagas
Why coordination is hard: the four impossibilities
- You cannot distinguish a slow node from a dead node. Every failure detector is a timeout, and every timeout is a guess. This is why fencing tokens exist.
- Clocks lie. NTP drift, leap seconds, and VM pauses mean wall-clock ordering across machines is unreliable. Use logical clocks (Lamport, vector clocks) for ordering; use wall clocks only for human-facing timestamps and TTLs.
- The FLP result — in a fully asynchronous system with even one faulty process, no deterministic algorithm guarantees consensus. Practical systems escape this with timeouts and randomisation.
- The two generals problem — no finite exchange of messages over a lossy channel gives both sides certainty. Hence: no true exactly-once delivery, only idempotency.
Consensus: Raft in one page
You will almost never implement consensus. What matters is knowing when you need it (exactly one leader, cluster membership, configuration, distributed locks, uniqueness) and that it costs a majority round trip on every write — which is why you keep the coordination layer small and off the hot path. Storing user posts in etcd is a wrong answer; storing "who is the current leader of shard 7" in etcd is right.
Distributed locks — and why they are a trap
Ask whether the operation can be made idempotent, atomic in a single store, or serialised through a partition. "Route all operations for a given key to the same partition/consumer, so they are processed one at a time" replaces a distributed lock with a data-layout decision, and it is faster and safer. Kafka partitioning, actor models, and per-entity single-threaded processing are all this idea.
Distributed transactions: 2PC vs Saga
| Two-phase commit | Saga | |
|---|---|---|
| How | Coordinator asks all participants to prepare; if all vote yes, it tells all to commit | A sequence of local transactions; each has a compensating transaction to undo it |
| Consistency | Atomic across services | Eventually consistent; intermediate states are visible |
| Failure mode | Coordinator dies after prepare → participants hold locks indefinitely ("in-doubt") | Compensation may itself fail; needs retries and manual reconciliation |
| Availability | Poor — any participant down blocks everyone | Good — each step is a local transaction |
| Verdict | Avoid across services. Fine inside one database. | The standard answer for microservices. |
Practice — Design "reserve a seat" across a Booking service and a Payment service without 2PC.distributed
1. Reserve with a TTL. Booking writes a HOLD row for the seat with a 10-minute expiry and a unique constraint on (showId, seatId) so only one hold can exist. This is a local transaction — atomic, no coordination.
2. Charge. Payment is called with an idempotency key derived from the hold ID, so retries cannot double-charge.
3. Confirm or compensate. On success, Booking converts HOLD → CONFIRMED. On failure or timeout, the hold simply expires and the seat returns to the pool — no compensation is needed at all, because the TTL is the compensation. If the charge succeeded but confirmation failed, a reconciliation job matches payments against bookings and either completes the booking or refunds.
Why this is the right answer: it converts a distributed transaction into a local transaction plus an expiring lease plus idempotent retry — three things that each work independently. And the "TTL as automatic compensation" observation is exactly the kind of insight interviewers are listening for.
Observability, deployment, and multi-region
The three pillars, and what each is for
- Metrics — cheap, aggregated numbers over time. Answer "is something wrong?" Use them for dashboards and alerts. Watch out for high-cardinality labels (never put userId in a Prometheus label — it will destroy the store).
- Logs — expensive, high-detail events. Answer "what exactly happened to this request?" Make them structured (JSON), sampled at high volume, and always carry a trace ID.
- Traces — the path of one request across services with per-span timing. Answer "which hop is slow?" This is the only tool that makes microservice latency debuggable.
An average latency of 100 ms is compatible with 5% of users waiting 3 seconds. Report p50, p95, p99, and p99.9. And remember tail amplification: if a page makes 10 parallel calls each with a 1% chance of being slow, roughly 10% of page loads are slow. This is why p99 of a dependency becomes p50 of a user experience — a great line to drop in an interview.
The four golden signals
Latency (split successful vs failed — fast errors hide in an average), Traffic (requests/sec), Errors (rate and class), Saturation (how full is the most constrained resource: CPU, memory, connection pool, queue depth). If you can only have four dashboards, have these.
SLI / SLO / error budget
An SLI is a measurement ("proportion of requests served in under 300 ms"). An SLO is the target ("99.9% over 28 days"). The error budget is the remaining 0.1% — and its purpose is to make reliability a negotiation rather than an argument: budget left means ship features, budget exhausted means freeze and fix. Saying this framing back to an interviewer shows you understand why SLOs exist organisationally, not just technically.
Deployment strategies
| Strategy | How | Trade-off |
|---|---|---|
| Rolling | Replace instances a few at a time | No extra capacity needed; both versions run simultaneously, so your API must be backwards-compatible |
| Blue-green | Two full environments; flip the LB | Instant rollback; costs 2× infrastructure during the switch; database migrations are the hard part |
| Canary | 1% → 5% → 25% → 100%, watching metrics | Best risk control; needs good metrics and automated rollback |
| Feature flags | Deploy dark, enable per cohort | Decouples deploy from release; enables instant kill-switch; flags accumulate as tech debt |
Use the expand–migrate–contract sequence: (1) expand — add the new nullable column/table, deploy code that writes both old and new; (2) migrate — backfill in batches, never one giant UPDATE that locks the table; (3) switch reads to the new column; (4) contract — stop writing the old one, then drop it in a later release. Never rename a column in a single deploy — during a rolling deploy, old and new code run at the same time. Volunteering this is a strong production-experience signal.
Multi-region
| Model | Writes | RTO / RPO | Complexity |
|---|---|---|---|
| Single region + backups | One region | Hours / minutes-to-hours | Lowest — and correct for most products |
| Active-passive | One region; async replica standby | Minutes / seconds | Moderate; must rehearse failover or it will not work |
| Active-active | All regions accept writes | Seconds / near-zero | High — you now own conflict resolution |
Active-active forces a choice on write conflicts: partition users by home region so a given record is only written in one place (best), use CRDTs for genuinely concurrent data, or accept last-write-wins and its data loss. Also state the two acronyms explicitly — RTO is how long recovery takes; RPO is how much data you can afford to lose.
Security, auth, and abuse control
Sessions vs tokens
| Server-side session | JWT (stateless token) | |
|---|---|---|
| Storage | Session store (Redis); cookie holds only an ID | Nothing server-side; the token carries signed claims |
| Revocation | Immediate — delete the session | Hard — the token is valid until it expires. Needs a denylist (which reintroduces state) or very short TTLs |
| Scale | One lookup per request (~0.5 ms, cached) | No lookup; verify a signature locally |
| Best for | First-party web apps where logout must be instant | Service-to-service auth, mobile, federated identity |
The production compromise: a short-lived access token (5–15 min, stateless, fast) plus a long-lived refresh token (stored server-side, revocable, rotated on each use with reuse detection). Say that pairing and the reasoning — you get stateless verification on the hot path and real revocation on the cold path.
OAuth 2.0 / OIDC in one paragraph
OAuth 2.0 is an authorisation delegation protocol: it lets an app act on a user's behalf at another service without the user handing over a password. OIDC layers authentication on top by adding an ID token. The flow you should name is authorisation code with PKCE: the client redirects to the identity provider, the user consents, the provider returns a one-time code to a registered redirect URI, and the client exchanges that code (plus a proof key it generated) for tokens. Note the common confusion in interviews: OAuth is not authentication — "log in with Google" is OIDC.
Abuse and rate limiting
Rate limiting protects capacity, cost, and fairness. The full design is case study H2; here are the algorithms in one table.
| Algorithm | Mechanism | Property |
|---|---|---|
| Fixed window | Count per minute bucket | Trivial. Allows 2× burst across a boundary (100 at 11:59:59, 100 at 12:00:00) |
| Sliding window log | Store every request timestamp (Redis ZSET), count those in the window | Exact. Memory grows with request count |
| Sliding window counter | Weighted blend of the previous and current window | Near-exact, O(1) memory. The usual production choice |
| Token bucket | Tokens refill at rate r, capacity b; a request takes one | Allows controlled bursts up to b. The most common API answer |
| Leaky bucket | Requests queue and drain at a fixed rate | Perfectly smooth output; adds queueing latency |
The rest of the security checklist
- Encryption — TLS 1.3 in transit; AES-256 at rest with keys in a KMS/HSM, rotated. Field-level encryption for the most sensitive columns.
- Passwords — bcrypt, scrypt, or Argon2id with a per-user salt. Never SHA-256, never MD5, never unsalted. Argon2id is the current recommendation.
- Authorisation — RBAC for coarse roles, ABAC for context-dependent rules. Check permissions at the service, never only in the UI. IDOR (changing
/orders/123to/orders/124) is the most common real-world API vulnerability. - Input — parameterised queries (SQL injection), output encoding + CSP (XSS), size limits and content-type validation on uploads.
- Secrets — a secrets manager, never environment variables committed to a repo; short-lived credentials via workload identity where possible.
- PII — data minimisation, retention limits, deletion pipelines that reach backups and derived stores. If you design for a global product, mention GDPR's right to erasure and how your event log complicates it (crypto-shredding: delete the per-user key so the ciphertext becomes unreadable — an elegant answer worth knowing).
Monoliths, microservices, and API design
The honest comparison
| Modular monolith | Microservices | |
|---|---|---|
| Deploy | One artifact; simple, atomic | Independent per team — the actual main benefit |
| Data | One database, real transactions, joins | Database per service; sagas, eventual consistency, no joins |
| Debugging | A stack trace | Distributed tracing across 8 services |
| Failure | All or nothing | Partial failure is the normal state, and you must design for it |
| Scaling | Scale the whole app | Scale only the hot service |
| Right when | < ~30 engineers, unclear domain boundaries, early product | Many teams needing independent release cadence; genuinely different scaling profiles |
Proposing microservices for a system with 100 users, or splitting services by technical layer ("a database service", "a caching service") rather than by business capability. The correct answer to "monolith or microservices?" almost always starts with "a well-modularised monolith, with clear internal boundaries so we can extract a service when a specific team or scaling need demands it" — and then naming the first service you would extract and why. Distributed systems are a cost you pay for organisational scale, not a badge.
API gateway — what belongs in it
Cross-cutting concerns that should not be reimplemented in every service: TLS termination, authentication and token validation, rate limiting, request routing and versioning, request/response transformation, caching of idempotent GETs, and observability injection (trace IDs). What does not belong: business logic. A gateway that knows about orders is a distributed monolith with extra latency.
A BFF (Backend for Frontend) is a per-client-type gateway — one for mobile, one for web — that aggregates and reshapes responses for that client's needs. It exists because a mobile client on a slow network wants one 8 KB call, not six 2 KB calls. Service mesh (Envoy sidecars, Istio/Linkerd) moves retries, mTLS, circuit breaking, and traffic-splitting out of application code and into the infrastructure layer — a legitimate answer when asked "how do you enforce retry policy consistently across 40 services?"
API design details interviewers probe
- Versioning — URI (
/v1/) is the pragmatic default; header-based is purer but harder to debug. Better still: never break — add fields, never remove or repurpose them, and treat unknown fields as ignorable. - Pagination — cursor over offset (F1). Return an opaque
nextCursorso you can change the underlying strategy later. - Idempotency — an
Idempotency-Keyheader on all unsafe operations, with the result stored for 24 hours. - Errors — a consistent envelope with a machine-readable code, a human message, and a request ID the user can quote to support. Use the right status: 400 malformed, 401 unauthenticated, 403 unauthorised, 404 missing, 409 conflict, 422 semantically invalid, 429 rate-limited, 503 overloaded.
- Bulk endpoints — one call that accepts 100 IDs is worth more than any latency micro-optimisation on a chatty client.
- Long-running work — return
202 Acceptedwith a status URL rather than holding the connection.
Practice — When would you deliberately choose a monolith at a company that already runs microservices?judgement
When the new system is (a) owned end-to-end by one team, (b) has strong transactional requirements across its entities, and (c) has an unclear or still-changing domain boundary. Splitting a domain you do not yet understand is the most expensive mistake available, because service boundaries are far harder to move than module boundaries — a refactor across two services requires a coordinated deploy, a data migration, and probably a saga.
The strong version of the answer names the exit criteria: "I'd build it as one deployable with hard internal module boundaries and no shared database tables between modules. The moment one module needs a different scaling profile or a second team owns it, extraction is a mechanical change rather than a rewrite." Interviewers are testing whether you optimise for reversibility.
Part II — twenty case studies
Each case study follows the same skeleton, which is also the skeleton you should follow in the room: requirements → numbers → API and data model → simple design → the one hard problem → failure modes → follow-ups. The "one hard problem" line is the important one — every classic question has a single central difficulty, and interviewers are waiting to see whether you find it.
| # | Question | The one hard problem |
|---|---|---|
| H1 | URL shortener | Generating short unique keys without coordination |
| H2 | Rate limiter | Making counters correct and cheap across many nodes |
| H3 | Unique ID generator | Sortable, unique IDs with no central bottleneck |
| H4 | Distributed key-value store | Quorums, replica placement, conflict resolution |
| H5 | Web crawler | Politeness + frontier prioritisation + dedupe at scale |
| H6 | Notification system | Fan-out, third-party unreliability, deduplication |
| H7 | News feed | Fan-out on write vs read, and the celebrity problem |
| H8 | Chat system | Connection routing, ordering, offline delivery |
| H9 | Search autocomplete | Sub-100 ms prefix lookup over a huge, changing corpus |
| H10 | YouTube | Transcoding pipeline and petabyte-scale delivery |
| H11 | Google Drive | Delta sync, chunking, conflict resolution |
| H12 | Uber | Geospatial indexing and real-time matching |
| H13 | Ticketmaster | Strong consistency and holds under extreme burst |
| H14 | Payment system | Idempotency, ledgers, reconciliation |
| H15 | Ad click aggregator | Exactly-once aggregation on a firehose |
| H16 | Distributed job scheduler | At-least-once execution at the right time |
| H17 | Metrics & monitoring | Cardinality, downsampling, write amplification |
| H18–H20 | Rapid fire: S3, Docs, leaderboard, wallet, stock exchange, proximity | Assorted |
Design a URL shortener (TinyURL / bit.ly)
Requirements
- Functional: shorten a long URL; redirect a short URL to the original; optional custom alias; optional expiry; basic click analytics.
- Non-functional: redirect latency < 100 ms at p99 (it sits in front of a page load); very high availability (a dead shortener breaks every link ever shared); read-heavy 10:1 or more; short links must be non-guessable enough to not enumerate a competitor's links.
- Scale: 100M new URLs/day, 1B redirects/day → ~1,000 writes/s and ~12,000 reads/s average, ~35,000 reads/s peak.
How short is short?
Key generation — the actual interview question
| Approach | How | Verdict |
|---|---|---|
| Hash the URL (MD5/SHA → base62, take first 7) | Deterministic; same URL → same key | Collisions must be detected and resolved with a retry/salt loop, which costs an extra read per write. Also leaks that two users shortened the same URL. |
| Auto-increment ID → base62 | DB sequence, encode the integer | Zero collisions, shortest keys. But keys are sequential and therefore enumerable (a crawler can walk your entire link database), and the sequence is a central bottleneck. |
| Random + collision check | Generate 7 random base62 chars, INSERT ... IF NOT EXISTS | Simple, unguessable. Collision probability is negligible until the space fills (birthday bound), and the DB's unique constraint handles it. |
| Pre-generated key store ✅ | A batch job pre-computes millions of unique random keys into an "available keys" table/queue; each app server leases a block of 10,000 into memory | Best answer. Key generation is off the write path entirely, writes never collide, no coordination per request, and keys are unguessable. Cost: a background service and the acceptance that a server crash wastes its unused block (harmless). |
"I'd use a counter-based approach for compactness but not expose the raw counter — I'd run the ID through a bijective permutation (e.g. a Feistel network or multiply by a large odd number mod 62⁷) before base62 encoding. That preserves uniqueness with zero collisions while making the output non-sequential." This is the most elegant answer and very few candidates give it.
Data model and API
Architecture
301 or 302? A 301 (permanent) lets the browser and intermediaries cache the redirect forever — which is great for latency and load, and terrible for analytics, because subsequent clicks never reach you. A 302 (temporary) guarantees every click hits your server. Answer: 302 if analytics matter, 301 with a short Cache-Control max-age if load matters. Being asked this and knowing the analytics implication is the entire point of the question.
Bottlenecks and follow-ups
- Hot links — one viral link can be 10% of all traffic. It will be cached, so this is mostly fine; add an in-process LRU on each app node for the top few thousand keys to avoid even the Redis hop.
- Analytics — never write to a database on the redirect path. Emit to Kafka, aggregate in a stream processor, store in a columnar store. Redirect latency must not depend on analytics availability.
- Expiry — a TTL on the KV store handles it for free; otherwise a batch cleanup job. Say that TTL cuts your 73 TB estimate dramatically.
- Abuse — shorteners are used to hide phishing. Check submissions against a safe-browsing API asynchronously and disable links on a hit; rate-limit creation per account and IP.
- Custom aliases — same table, but writes must go through a uniqueness check, and you should reserve a blocklist (profanity,
admin,login).
Follow-up — Two users shorten the same long URL. Same short key or different?product
Different, in almost every real product. Reusing the key looks like a storage optimisation, but it breaks per-user analytics (whose clicks are these?), per-user expiry, and revocation (deleting one user's link would break the other's). It also leaks information across accounts. The storage "saved" is a few hundred bytes per duplicate — irrelevant next to the product problems. The correct answer names the trade-off and picks the product side, which is the behaviour being tested.
Follow-up — How do you make this globally fast?scale
Redirects are read-only and immutable once created, which makes them nearly perfect CDN/edge material. Serve them from edge workers with a regional cache; on a miss, read from a regional read-replica of the KV store. Writes can stay in one home region since 1,000 writes/s is trivial and creation latency is far less sensitive than redirect latency. State the consistency consequence honestly: a newly created link may take a second or two to become resolvable at every edge, which is acceptable — and if it is not, the creation response can pre-warm the edge cache for the creator's own region.
Design a distributed rate limiter
Requirements
- Functional: limit requests per identity (user, API key, IP) per rule; return 429 with
Retry-AfterandX-RateLimit-*headers; support multiple tiers (free/paid) and multiple granularities (per-second and per-day simultaneously). - Non-functional: adds < 5 ms to a request; must not become a single point of failure; must work across hundreds of API server nodes; approximate is acceptable at the edges.
Where does it live?
Three options: client-side (unreliable — a client can simply not implement it), in the API gateway / middleware (the usual answer — centralised, protects everything behind it), or as a separate service (flexible, but adds a network hop to every request). Choose the gateway, and say why: it is the first place that can reject a request before it consumes any downstream capacity.
Algorithm choice
See the table in F18 for all five. In a design round, pick token bucket for public APIs (bursts are user-friendly and it is trivially explained) or sliding window counter when smoothness matters more than burst tolerance. Then show the implementation.
Making it distributed
The naive approach — each node keeps its own counter — means N nodes allow N× the limit. Three fixes:
- Central store (Redis). All nodes read/modify/write the same key. Correct, but the read-modify-write must be atomic or concurrent requests will over-admit. Use a Lua script (single round trip, atomic on Redis's single-threaded executor). This is the standard answer.
- Sharded limits. Give each of N nodes limit/N and route each identity consistently to one node (consistent hashing). No shared state and no network hop, but uneven traffic wastes budget.
- Local + async sync. Each node enforces locally and gossips counts every few hundred ms. Fast and highly available, over-admits slightly during the sync interval. Used by systems that need microsecond decisions.
(1) The TTL. Without PEXPIRE, you store a key forever for every user who ever made one request. With it, memory is proportional to active users. (2) Fail open or fail closed? If Redis is down, do you block all traffic or allow it? For a public API protecting a fragile backend, fail closed with a local fallback limiter; for a user-facing product where availability matters more than precision, fail open and alert. Stating that this is a deliberate product choice is exactly the judgement being graded.
Response contract
Follow-up — How would you rate-limit by IP when users are behind NAT/CGNAT?nuance
IP is a poor identity: an entire university or mobile carrier can share one address, so a per-IP limit either punishes thousands of legitimate users or is set so high it stops nothing. Layered answer: (1) prefer authenticated identity (API key, userId) whenever available, and use IP only for unauthenticated endpoints like login and signup; (2) for those, use a much higher IP limit combined with a strict per-account limit, so credential stuffing across many accounts is caught by the IP rule and brute force on one account by the account rule; (3) add a proof-of-work or CAPTCHA step rather than a hard block, so shared-IP users have a path forward; (4) if you sit behind a proxy or CDN, read the real client IP from a trusted header and never from an untrusted X-Forwarded-For a client can spoof.
Design a unique ID generator (Snowflake)
Requirements
Globally unique; 64-bit (fits a BIGINT and a Java long); roughly time-sortable (so that a primary-key index has good insert locality and "newest first" needs no extra index); at least 10,000 IDs/sec per node; no central coordinator on the hot path.
| Approach | Problem |
|---|---|
| DB auto-increment | Single point of failure and a write bottleneck; cannot span shards |
| UUIDv4 (random 128-bit) | Unique and coordination-free, but 128 bits, not sortable, and random inserts destroy B-tree locality → page splits and write amplification |
Ticket server (a single MySQL doing REPLACE INTO) | Simple, used at Flickr, but a SPOF and a network round trip per ID |
| Snowflake ✅ | Coordination-free after startup, sortable, 64-bit. Depends on clocks and on unique machine IDs. |
| UUIDv7 / ULID | Modern alternative: 128-bit but time-prefixed, so it sorts and indexes well. Mention it — it is what many teams pick today when 128 bits is acceptable. |
"Where does machineId come from?" — from ZooKeeper/etcd as an ephemeral sequential node, from a Kubernetes StatefulSet ordinal, or derived from the pod IP. Whatever you choose, you must guarantee no two live instances share one, and you must handle re-use after a pod restarts. "What if the clock goes backwards?" — refuse to issue IDs until time catches up (as above) and alert; issuing duplicates is far worse than a brief 503. Also disable NTP slew-jumps on ID nodes.
Follow-up — Are Snowflake IDs strictly ordered across the cluster?precision
No — only roughly. Two IDs generated in the same millisecond on different machines are ordered by machineId, not by real time, and clock skew of a few milliseconds between nodes means an ID generated slightly later can sort earlier. This is fine for "newest first" feeds and for index locality, and it is not fine if you need a total order for correctness (e.g. deciding which of two conflicting writes wins). For that you need consensus, a single sequencer, or logical clocks. Volunteering this distinction — approximate ordering for indexing vs strict ordering for correctness — is a strong signal.
Design a distributed key-value store
This is really an oral exam on F8, F9, and F10 combined. The expected answer is a Dynamo-style design; walk through it as a set of decisions.
Requirements
get(key) and put(key, value) with values up to ~10 KB; scale to petabytes; configurable consistency; high availability during node and network failures; automatic scaling and rebalancing.
The seven decisions
| Decision | Choice | Because |
|---|---|---|
| Data partitioning | Consistent hashing with virtual nodes (F9) | Adding/removing a node moves only 1/N of keys; even distribution |
| Replication | N=3, the next 3 distinct physical nodes clockwise, placed in different racks/AZs | Survives a rack or AZ failure, not just a machine failure |
| Consistency | Tunable quorum: W + R > N (F10) | The caller picks per-request: fast-and-eventual or slow-and-consistent |
| Conflict resolution | Vector clocks → sibling values returned to the client; LWW only where the app can tolerate loss | Detects genuine concurrency instead of silently discarding a write |
| Failure detection | Gossip protocol with heartbeat counters | Decentralised; no coordinator to lose. Every node eventually learns membership |
| Temporary failure | Hinted handoff — a healthy node accepts the write on behalf of the down node and replays it when it returns | Keeps write availability at W even during a failure ("sloppy quorum") |
| Permanent failure | Anti-entropy with Merkle trees | Compares replicas by hashing ranges recursively, so only the differing ranges are transferred, not the whole dataset |
Follow-up — A client writes with W=1 and immediately reads with R=1. What can go wrong?consistency
W + R = 2 which is not > N = 3, so the read set and write set need not overlap: the read can land on two replicas that have not yet received the write, and return the old value or nothing at all. The user sees their own write disappear.
Fixes in order of cost: (1) use W=2, R=2 for that key class — a quorum guarantees overlap; (2) keep W=1 but route the client's reads to the coordinator node that accepted the write for a short session window (read-your-writes, F5); (3) have the client cache what it just wrote and merge locally, which is what most mobile apps do anyway.
The framing that matters: W=1/R=1 is not a bug — it is a deliberate trade for latency and availability. The design error would be choosing it for data where a lost read is a correctness problem, like an account balance.
Design a web crawler
Requirements
- Functional: given seed URLs, download pages, extract links, repeat; store page content for downstream indexing; re-crawl pages based on how often they change.
- Non-functional: politeness (never overwhelm a host, obey
robots.txt), scale to a billion pages/month, extensible to new content types, robust against traps. - Scale: 1B pages/month ≈ 400 pages/s average, ~1,200/s peak; 500 KB average page → 500 TB/month of raw content.
The URL frontier is the whole problem
Deduplication at scale
- URL seen? — you cannot keep 100 billion URLs in a hash set. Use a Bloom filter (probabilistic, no false negatives, ~1% false positive means you occasionally skip a page — acceptable) backed by a sharded key-value store for exactness where it matters. Normalise URLs first: lowercase host, strip fragments, sort query parameters, drop tracking parameters.
- Content seen? — exact duplicates via a hash of the body. Near-duplicates (same article on 50 mirrors) via SimHash: documents within a small Hamming distance are treated as duplicates. Mentioning SimHash or MinHash specifically is the depth signal here.
Traps and hazards to volunteer
- Spider traps — infinitely deep calendar URLs, session-ID query parameters generating unbounded URLs. Defences: max depth, max URLs per domain, URL pattern blocklists, and manual review of the top-N domains by URL count.
- Robots.txt — fetch and cache per host with a TTL; respect
Crawl-delayandDisallow. Ignoring it gets you IP-banned and is a legal risk — say this. - JavaScript-rendered pages — a plain HTTP fetch gets an empty shell. A headless-browser rendering tier is 10–100× more expensive per page, so run it selectively based on whether the static HTML looked empty.
- Re-crawl policy — do not crawl everything at the same frequency. Estimate a per-page change rate from history and schedule accordingly; a news homepage may be hourly, a 2011 forum thread yearly.
- Coordination — shard the frontier by
hash(host)so politeness state for a host lives on exactly one node, which turns a distributed-locking problem into a partitioning decision.
Follow-up — A worker crashes holding 10,000 URLs it had dequeued. What happens?reliability
Nothing should be lost, and the standard fix is visibility-timeout semantics rather than a true dequeue: a worker leases URLs for, say, 5 minutes; if it does not acknowledge completion, the lease expires and the URLs become visible again for another worker. That gives at-least-once processing, which is fine here because crawling the same URL twice is harmless (the content-seen check absorbs it) — a nice example of designing so that duplicates do not matter instead of trying to prevent them.
Add: checkpoint progress in small batches rather than one big transaction; make the politeness delay table durable (or at worst re-derive it, since a fresh worker simply waits the default delay); and monitor lease-expiry rate as a signal that workers are dying or are too slow.
Design a notification system
Requirements
- Functional: send push (iOS/Android), SMS, email, and in-app notifications; triggered by services or by scheduled campaigns; users can opt out per channel and per category; templates with variables; retries.
- Non-functional: 10M notifications/day with bursts of 100k/s during campaigns; at-least-once delivery with deduplication; a third-party outage must not lose messages; soft real-time (seconds) for transactional, best-effort for marketing.
The five hard parts
- Device token management. A user has many devices; tokens expire, get revoked, or transfer to another user when a phone is resold. Store
(userId, deviceToken, platform, appVersion, lastSeen), and process the provider's feedback — APNs/FCM tell you when a token is invalid, and continuing to send to dead tokens gets your sender reputation throttled. - Deduplication. At-least-once queues mean a worker crash after sending but before acking causes a duplicate push, which users experience as spam. Fix with an idempotency key
(userId, notificationId, channel)written to a store with a TTL, checked before send (F12). - Third-party unreliability. Providers rate-limit you, go down, and return ambiguous errors. Each worker needs a circuit breaker, a token-bucket limiter matched to the provider's quota, exponential backoff, and a DLQ. For critical channels, keep a secondary provider and fail over — then be explicit that failover risks duplicates, so idempotency must be provider-independent.
- Priority and throttling. A marketing blast of 5M messages must never delay a password-reset email. Separate priority lanes (or separate topics) with dedicated worker pools, plus a per-user frequency cap ("no more than 3 marketing pushes per day") enforced centrally rather than by each producing service.
- Fan-out. "Notify all followers of X" can be one request producing 10M messages. Do not expand it in the API handler. Write a single campaign record, then have a fan-out job stream the recipient list in pages and emit per-recipient messages, with checkpointing so a crash resumes rather than restarts.
"Opt-out state is checked at send time, not at enqueue time, because a user may unsubscribe between the campaign being scheduled and it being delivered — and for SMS and email that is a legal requirement, not a nicety." Interviewers rarely expect it and always credit it.
Follow-up — How do you handle a 100k/s campaign burst without melting the SMS provider?flow control
The queue absorbs the burst — that is its job — and the workers drain it at a rate you control rather than at the rate the producer generated it. Concretely: a token-bucket limiter shared across workers (Redis, as in H2) set to slightly below the provider's contractual rate; a fixed worker pool sized to that rate rather than autoscaled on queue depth (autoscaling on queue depth is a classic mistake here — it scales you straight into the provider's rate limit and gets you throttled or blocked); and a per-campaign scheduler that spreads sends over a window ("deliver over 30 minutes") which is also better for user experience and for engagement metrics. Then monitor queue depth and estimated drain time, and alert if a transactional message's expected latency exceeds its SLO.
Design a news feed
Requirements
- Functional: post; follow/unfollow; view a home timeline of posts from people you follow, newest first (or ranked); view a user's own profile timeline.
- Non-functional: home timeline loads in < 200 ms at p99; extremely read-heavy (~1000:1); a few seconds of staleness is acceptable; high availability.
- Scale: 300M DAU, 300M posts/day (~3,500/s), ~300k timeline reads/s peak, average 200 followers with a long tail up to 100M+.
The core decision: fan-out on write vs fan-out on read
"Fan-out on write moves cost from the read path to the write path, which is correct when reads outnumber writes 1000:1. It breaks only for the celebrity tail, so I special-case that tail and merge at read time — the hybrid keeps read cost O(1 + small constant) while keeping write cost bounded."
Data model
Follow-ups you should pre-empt
- Ranking instead of chronological — the feed becomes: candidate generation (the merged set above, plus recommendations) → feature extraction → an ML ranker → business rules (diversity, ads, freshness). Say that ranking sits after candidate generation, so the storage design does not change.
- Unfollow / delete — do not attempt to remove the post from millions of materialised timelines. Filter at read time against a small "deleted posts" bloom filter or check on hydration, and let the entry age out of the capped list.
- New follow — backfill that user's recent posts into your timeline asynchronously, or simply let them appear from the next post onwards. Say which and why.
- Media — post rows carry only media IDs; the bytes are in object storage behind a CDN (F13).
Follow-up — Where exactly do you draw the celebrity threshold, and what happens at the boundary?nuance
The threshold is where fan-out cost exceeds merge cost, and it is empirical rather than principled — commonly somewhere in the tens of thousands of followers. The important part of the answer is what happens at the boundary: it must be a hysteresis band, not a single number, or an account hovering at the line flip-flops between strategies and produces duplicate or missing entries. Practical approach: promote to "celebrity" at 50k followers, demote only below 30k, and on transition either backfill or accept a brief window where both paths run and the read merges deduplicate by postId. Also make the flag a property of the author read at fan-out time, so it is a single lookup rather than a per-follower decision.
Design a chat system
Requirements
- Functional: 1:1 and group chat (up to 500 members); message delivery with sent/delivered/read receipts; online presence; offline message delivery; media attachments; message history.
- Non-functional: end-to-end latency < 200 ms for online users; no message loss; correct ordering within a conversation; 50M DAU with ~10M concurrent connections.
Message flow, step by step
- A sends over its WebSocket with a client-generated
clientMessageId(UUID). This is the idempotency key that makes retries safe when A's network flaps. - The gateway hands it to the message service, which assigns a server-side ID and, critically, a per-conversation sequence number.
- Persist before acknowledging. The ack to A means "durably stored", which is what a single tick means in the UI.
- Look up B in the session store. If online, publish to B's gateway node (via Redis pub/sub or Kafka keyed by nodeId) which writes it to B's socket → two ticks. If offline, the message simply stays in the store, and a push notification is sent.
- When B connects, it sends its last-seen sequence number per conversation and the server streams everything after it — a pull-based sync that is far more robust than trying to queue per-recipient.
The four hard parts
| Problem | Answer |
|---|---|
| Ordering | Do not rely on wall-clock timestamps from clients — phones have wrong clocks and networks reorder. Assign a monotonic per-chatId sequence number server-side; clients sort by it. Ordering only needs to be consistent within a conversation, which makes it cheap. |
| Routing across gateways | Session registry userId → {nodeId, connectionId} in Redis with a heartbeat TTL. Publish to the target node. If the lookup is stale (B reconnected elsewhere), the message falls back to the store and B pulls it on sync — so a stale registry causes latency, not loss. |
| Group chat fan-out | Store the message once per chatId, not once per member. Each member tracks their own read cursor. Fan-out is only for delivery, not storage. For 500-member groups this is trivially fine; the same pattern would not survive a million-member channel, which is why broadcast channels use a different (pull-based) model. |
| Storage choice | Cassandra partitioned by chatId, clustered by seq DESC. That gives "load the last 50 messages of this conversation" as a single-partition sequential read — the exact query the product makes 99% of the time. Cap partition size by bucketing very large chats by time. |
With E2EE (Signal protocol), the server stores and routes ciphertext it cannot read. Consequences you should name: server-side search is impossible (search must run on-device), group membership changes require re-keying, multi-device support requires per-device session keys and a key-distribution mechanism, and message history on a new device cannot simply be downloaded — it must be transferred from an existing device or restored from an encrypted backup. The architecture barely changes; the features change a lot. That framing is the answer.
Follow-up — How does presence ("last seen 2 minutes ago") work without melting the system?real-time
Presence is deceptively expensive because it is O(users × friends) updates. Three techniques: (1) heartbeat with TTL — the client sends a heartbeat every ~30 s, stored in Redis with a 60 s TTL; absence of the key means offline, so no explicit "went offline" event is needed and a crashed client is handled automatically; (2) do not broadcast presence to all contacts — fetch it on demand when a user opens a chat list, and subscribe to live updates only for the conversations currently on screen; (3) debounce and coarsen — "last seen" at minute granularity rather than second, so a flapping mobile connection does not generate a storm of updates. Add the privacy angle: presence is a setting, and many users disable it, which conveniently also reduces load.
Design search autocomplete (typeahead)
Requirements
- Functional: given a prefix, return the top 5–10 most popular completions; support new/trending queries; ignore inappropriate terms.
- Non-functional: < 100 ms end-to-end — it must feel instantaneous while typing; extremely read-heavy; results may be minutes stale.
- Scale: 10M searches/day; every keystroke is a potential request, so ~20 requests per search → a 20× amplification over search traffic itself.
The data structure: trie with precomputed top-k
Architecture: separate the read path from the build path
Scaling the trie
A full trie over hundreds of millions of terms will not fit in one process. Shard by prefix — but not naively by first letter, because 's' has vastly more terms than 'z'. Use a learned prefix split: measure the distribution and assign ranges so each shard carries similar load (e.g. shard 1 = 'a'–'ac', shard 2 = 'ad'–'am'…), with a small routing map cached on the client-facing tier. Since the trie is rebuilt offline anyway, recomputing the split each build is free.
"Most of the 100 ms budget is network, not computation. So: debounce keystrokes on the client, cancel in-flight requests when a new keystroke arrives, cache aggressively in the browser, serve from an edge PoP, and keep the whole data structure in memory so the server side is a sub-millisecond map lookup. I'd also return results for prefixes the user hasn't typed yet — prefetching the next character's results is cheap and makes it feel instant."
Follow-up — How do you personalise it without breaking the cache?nuance
Personalisation and caching are in direct tension: a per-user response cannot be shared at the edge. The standard resolution is a two-layer merge. The global top-k stays cacheable and is served from the edge exactly as before. The personal layer — your own recent searches, your location's trending terms — is small, and can either live entirely on the client (the browser stores your last 100 queries and merges them into the displayed list locally, which costs zero server latency and is also better for privacy) or come from a small per-user list fetched once per session rather than per keystroke. The merge happens at render time. This keeps 100% edge cacheability for the expensive part while still feeling personalised — and stating that trade explicitly is what the question is testing.
Design YouTube
Requirements
- Functional: upload video; watch video with adaptive quality; view counts; basic search and recommendations (park these).
- Non-functional: playback starts in < 1 s; no buffering on variable networks; upload of a 10 GB file must survive a dropped connection; global audience.
- Scale: 500 hours uploaded per minute; 5M concurrent viewers; average watch bitrate 2 Mbps → 10 Tbps of egress. That number, not storage, is the design driver — and cost is a legitimate first-class concern here.
Upload path
Transcoding: a DAG, not a script
Playback: adaptive bitrate streaming
The video is stored as many independent segments per quality level, described by a manifest (HLS .m3u8 or DASH .mpd). The player — not the server — decides which quality to request for the next segment based on measured throughput and buffer level. That is why streaming scales: the server just serves static files, and all the adaptation logic is on the client. Start with a low-bitrate segment so playback begins fast, then ramp up.
Delivery and cost
- CDN is mandatory, and the cache hit ratio is the whole economic story. Popularity is extremely skewed — a small fraction of videos gets most views — so a modest edge cache serves the vast majority of traffic. Long-tail videos are served from a regional origin.
- Tiered storage — hot videos on fast storage with all renditions; cold videos moved to cheap storage, and rarely-watched renditions deleted and regenerated on demand. Naming this trade (storage cost vs re-transcode cost) is a strong signal.
- Pre-positioning — for a scheduled premiere or a major sports event, push content to edge caches before the event rather than letting 5M simultaneous first-viewers all miss.
- Signed URLs with short expiry for access control, so the CDN can serve private content without calling your auth service per segment.
- View counts are not a database
UPDATE— they are events into a stream aggregator (H15), with dedupe rules (a "view" is typically 30 s of playback, once per user per window, to resist inflation).
Follow-up — A video goes viral and 2M people start watching in 60 seconds. What breaks?scale
Not the origin, if the CDN is doing its job — but the first minute is the danger, because before the segments are cached at each edge, the misses all arrive at the origin simultaneously (a stampede, F11). Defences: request coalescing at the edge (many CDNs support this natively — 1,000 concurrent misses for the same object become one origin fetch), a mid-tier shield cache between edges and origin so the origin sees one request per region rather than per PoP, and pre-warming when a video's view velocity crosses a threshold.
What else breaks: the metadata service (fix: cache the video metadata aggressively, it is immutable once READY), the view-count path (fix: it is already async and sampled), and the comments/live-chat system if there is one (a genuinely different problem — that one needs its own fan-out design). The instinct to check every component against the spike, rather than only the obvious one, is what is being graded.
Design Google Drive / Dropbox
Requirements
- Functional: upload/download files; sync across devices automatically; share with other users; version history; work offline and reconcile on reconnect.
- Non-functional: sync must be bandwidth-efficient (do not re-upload a 2 GB file because one byte changed); strong durability; conflicts must never silently lose data.
The central idea: chunking + content-addressed storage
Sync protocol
- Client watches the local folder and computes chunk hashes for changed files.
- It asks the server which of those chunk hashes are already known (
POST /chunks/check) — a cheap way to skip uploading anything the system already has. - It uploads the missing chunks directly to object storage, then commits a new file version to the metadata service with the previous version ID (an optimistic concurrency check).
- Other devices are told "namespace changed since cursor X"; they pull the metadata delta and download only the chunks they lack. A cursor-based delta, not a full listing — this is what keeps sync O(changes) instead of O(files).
Conflicts
Two devices edit offline and both commit against version 7. The second commit fails its optimistic check. Do not merge binary files and do not overwrite. The correct product behaviour is to keep both: create report (Alice's conflicted copy).pdf and surface it. For text-like formats, an operational-transform or CRDT layer can merge (see H19), but for general files, preserving both and telling the user is the honest answer — and saying "silently losing a user's work is the one unacceptable outcome" frames it correctly.
Follow-up — A user shares a folder with 10,000 people. What is the sync problem?fan-out
Every change to that folder must be reflected in 10,000 users' change feeds — a fan-out problem structurally identical to the celebrity case in H7. The fix is the same shape: do not materialise the change into 10,000 feeds. Model the shared folder as its own namespace with its own monotonically increasing cursor; a user's client subscribes to the namespaces it has access to and polls each cursor. The change is stored once, and each of the 10,000 clients discovers it by comparing cursors. Permissions are then also evaluated per namespace rather than per file, which keeps ACL checks O(1) instead of O(files).
The transferable insight — when fan-out is expensive, introduce a shared object with a version cursor and let readers pull — recurs in feeds, chat channels, and collaborative documents.
Design Uber (proximity + matching)
Requirements
- Functional: drivers broadcast location; riders request a ride; the system matches a nearby driver; both track each other live; trip lifecycle and pricing.
- Non-functional: matching in < 5 s; location updates every 4 s from every active driver; a rider must never be matched to two drivers, nor a driver to two riders.
- Scale: 1M active drivers → 250k location writes/s. That write rate is the first thing to notice, and it rules out a normal relational table with a spatial index.
Geospatial indexing — the actual question
| Technique | How it works | Trade-off |
|---|---|---|
| 2D lat/lng index | B-tree on each column | Only one dimension is usable per query; a radius search degenerates into a bounding-box scan. Poor. |
| Geohash | Interleave lat/lng bits into a base32 string; a shared prefix means spatial proximity | Simple, string-prefix-searchable in any KV store, and it is what Redis GEO uses. Edge case: two nearby points can differ at the first character (across a cell boundary), so you must search the cell and its 8 neighbours. |
| Quadtree | Recursively split space into 4 until a cell holds < k points | Adapts to density — dense Manhattan cells subdivide, empty ocean cells do not. In-memory structure, rebuilt/updated dynamically. |
| S2 / H3 | Project the sphere onto a cube (S2) or hexagonal grid (H3), with a hierarchical cell ID | The production answer. Hexagons (H3, used by Uber) have uniform neighbour distance, which matters for dispatch. Cell IDs are integers, so they index and shard beautifully. |
Geo-sharding creates a natural hotspot problem: a stadium at closing time puts enormous load on one cell. Two answers: (1) cell size should adapt to density, which is exactly what quadtrees and multi-resolution S2/H3 give you; (2) surge pricing is not just a business feature — it is a load-shedding mechanism that reduces demand in an overloaded supply region. Framing surge as demand shaping rather than revenue is a memorable answer.
Follow-up — How do you show the rider a live-moving car without 250k writes/s hitting your API?real-time
Only the driver on an active trip needs to stream to a specific rider, so the fan-out is 1:1, not broadcast — that is already a huge reduction. The driver's app sends location over an existing WebSocket to a gateway; the gateway looks up the trip and forwards to the rider's connection directly, without a database write in the path. Reduce further by: sending updates every 4–5 s rather than continuously; having the client interpolate smoothly between points (the smooth animation is a client-side illusion, not a data-rate requirement); snapping to the road network so the car does not appear in buildings; and compressing the payload to a few bytes of delta-encoded coordinates. Persist the route asynchronously in batches for receipts and disputes, decoupled from the live path entirely.
Design a ticket booking system
This is the counterpoint to the feed questions: here eventual consistency is wrong. Double-selling one seat is a business failure, so the design must be consistent where it counts and only optimistic elsewhere.
Requirements
- Functional: browse events; see a seat map; hold seats while paying; complete purchase; release abandoned holds.
- Non-functional: no double-booking, ever; survive a burst of 100k users at an on-sale moment for 20k seats; the browse path should stay available even under that burst.
Design
"The scale here is misleading. 100k concurrent users sounds like a distributed-systems problem, but there are only 20k seats — the contended state is tiny. So I'd keep all seat state for one event on a single node with real transactions, and spend the scaling effort on the read path and on the queue in front of it. Distributing the thing that needs consistency would be solving the wrong problem."
Handling the on-sale burst
- Virtual waiting room. Admit users into the purchase flow at a controlled rate (a token/queue system at the edge). This converts an unbounded burst into a steady stream the transactional core can absorb, and gives users an honest position indicator instead of an error page.
- Bot defence. The real burst is scalpers, not fans: CAPTCHA, device fingerprinting, per-account and per-payment-instrument limits, and account-age requirements. Say it, because it is the actual production problem.
- Hold TTL tuning. Too short and users on slow connections lose seats mid-payment; too long and inventory is locked by abandoned carts. 10 minutes with a visible countdown is the norm.
- Fairness. Randomised admission from the waiting room is fairer than strict FIFO (which rewards whoever had the fastest connection at the exact millisecond) — and less exploitable.
Follow-up — Would you use Redis to hold seats instead of the database?judgement
It is tempting — Redis is fast and atomic — but be careful about what you are trusting. Redis's default persistence is asynchronous, and a failover can lose the last few seconds of writes, which here means two people holding the same seat. If you use Redis, it must be as a front-line contention filter only, with the database remaining the source of truth: an atomic Redis operation quickly rejects the 99 users who lost the race, and only the winner proceeds to a database transaction that authoritatively confirms the hold. That gives you Redis's throughput and the database's durability, and it degrades safely — if Redis is down you can fall back to hitting the database directly, just slower.
Answering "yes, Redis, it's atomic" without the durability caveat is the trap. Recognising that atomic and durable are different properties is the point of the question.
Design a payment system
Requirements
- Functional: accept a payment for an order; support cards via a PSP (Stripe/Adyen); handle refunds; produce an auditable record of every money movement; reconcile with the provider daily.
- Non-functional: never double-charge; never lose a payment record; every state change auditable and immutable; correctness beats latency and beats availability.
Three non-negotiable design elements
You call the PSP and the connection times out. Did the charge happen? You do not know — and this is the defining problem of payment engineering. Never retry blindly. The correct sequence: record your intent before calling out (so a crash leaves evidence), retry with the same idempotency key so the PSP deduplicates, and if uncertainty persists, query the PSP for the status of that key rather than issuing a new charge. Design so that "unknown" is a legal state with a resolution path, rather than something you hope never happens.
"I'd never let raw card numbers touch my servers — the client tokenises directly with the PSP's SDK and I store only the token. That reduces PCI-DSS scope from the whole system to almost nothing." One sentence, and it demonstrates you know what actually makes payment systems hard to build.
Design a real-time analytics / ad click aggregator
Requirements
- Functional: ingest click events; return click counts per ad per minute; support "top N ads in the last M minutes"; support filtering by dimensions (country, device).
- Non-functional: 1M events/s; dashboard freshness within ~1 minute; accuracy matters — advertisers are billed from these numbers, so duplicates and losses are financial errors; must tolerate late-arriving events.
The four hard parts
- Exactly-once counting. At-least-once delivery plus naive counting means over-billing. Solutions: assign each click a unique
clickIdat the edge and deduplicate within the window (a Bloom filter or a per-window key set); use Flink's checkpointed state with two-phase-commit sinks; or make the sink idempotent by writing(adId, minute) → countas an overwrite of a recomputed window value rather than an increment. The overwrite trick is the simplest to explain and defend. - Event time vs processing time. A click at 12:00:59 may arrive at 12:01:30 because of a mobile network delay. Aggregating by arrival time puts it in the wrong minute. Use event time with a watermark — "I will wait up to 2 minutes for stragglers, then close the window" — and route anything later than that to a late-data path that issues a correction. Being able to explain watermarks concretely is the depth signal in this question.
- Hot keys. A viral ad concentrates on one Kafka partition. Fix with two-stage aggregation: pre-aggregate locally per task with a salted key (
adId#0..9), then combine — the same counter-sharding idea from F10. - Fraud. Click fraud is the real product problem: dedupe by
(userId, adId, minute), filter known bot signatures, cap per-IP rates, and keep a separate "billable" pipeline whose output is what advertisers are charged, downstream of fraud filtering, so the raw and billable numbers can differ and be reconciled.
Follow-up — The dashboard shows 4.2M clicks; the billing system says 4.1M. What do you do?correctness
First, do not "fix" either number until you know why they differ — this is a reconciliation problem, exactly as in H14. Enumerate the legitimate causes before assuming a bug: fraud filtering applied to billable but not to the dashboard; different window boundaries or time zones; late events that landed after the dashboard snapshot but before billing closed; duplicates removed in one path but not the other; and the possibility that the dashboard reads a cached, stale aggregate.
The systemic answer is to make the discrepancy measurable and expected rather than surprising: emit both counts through the same lineage with an explicit filtered_reason breakdown, so the difference decomposes into named categories that sum exactly. Then set an alert on any unexplained residual. The general principle — reconcile by decomposition, not by guessing — is what transfers to every data-correctness question.
Design a distributed job scheduler (cron at scale)
Requirements
Schedule one-off and recurring jobs; execute at (or shortly after) the requested time; at-least-once execution with idempotent handlers; retries with backoff; visibility into job state; millions of scheduled jobs, thousands executing per second.
"What if a worker dies mid-job?" The lease expires and another worker picks it up — hence at-least-once, hence handlers must be idempotent. "What if the scheduler is down for 10 minutes?" On recovery, thousands of overdue jobs fire at once (a thundering herd). Mitigate by rate-limiting the poller's dispatch, and by giving each job a policy: some must run late, others should be skipped if their window has passed (a stale "send morning digest" at 3 pm is worse than not sending it). "Exactly-once?" No — at-least-once plus idempotency, as always.
Design a metrics and monitoring system
Requirements
Collect metrics from 100k hosts; store at 10-second resolution; query arbitrary time ranges with aggregation; alert on thresholds; retain 13 months. That is ~10M data points/second at peak — the write path is the design.
Evaluate rules on a schedule against the query engine; require a condition to hold for a duration (for: 5m) to suppress flapping; group and deduplicate alerts so one incident does not page you 400 times; support silences during deploys; and route by severity. Also: alert on symptoms (users are getting errors) rather than causes (CPU is 90%) — CPU at 90% with happy users is not an incident, and this framing is straight out of production practice.
Rapid fire: six more systems
Shorter treatments. For each, the requirement that matters, the central insight, and the trap.
Design S3 / an object store
Central insight: separate the metadata plane (bucket → object key → chunk locations, a sharded KV store) from the data plane (immutable chunks on dumb storage nodes). Objects are immutable — an "update" writes a new version and repoints metadata, which makes replication and caching trivial. Durability comes from erasure coding across failure domains (F13); a background scrubber continuously reads and re-verifies checksums to catch bit rot before it compounds. Trap: proposing a POSIX filesystem. Object stores are flat namespaces precisely because directory trees require expensive hierarchical locking and rename semantics.
Design a leaderboard
Central insight: Redis sorted sets. ZADD leaderboard score userId, ZREVRANGE 0 9 for the top 10, ZREVRANK for a player's own rank — all O(log n). Scale problem: 50M players in one ZSET is memory-heavy and every write touches one key. Fix by sharding into buckets by score range and maintaining a small top-N set separately; approximate ranks for players outside the top few thousand ("top 12%") since nobody in rank 4,000,000 needs an exact number. Trap: computing rank with SELECT COUNT(*) WHERE score > x in a relational DB — correct and completely unscalable.
Design a digital wallet
Central insight: it is H14's ledger with a stricter constraint — balances cannot go negative, and transfers are between two internal accounts. Model every transfer as a double-entry pair inside a single transaction if both accounts are in one shard; if they are on different shards, use a saga with a reserved/pending state (F16). Trap: using a mutable balance column. Use the ledger plus periodic snapshots, and enforce the non-negative constraint at the database level (CHECK or a conditional update), never only in application code.
Design Google Docs (collaborative editing)
Central insight: concurrent edits to shared text need either OT (operational transformation — transform incoming operations against concurrent ones so all clients converge; what Google Docs uses, requires a central server to order operations) or CRDTs (data types whose merge is commutative and associative by construction, so no central authority is needed; used by Figma, Yjs, Automerge). Both converge; OT is more bandwidth-efficient but harder to get right, CRDTs are simpler to reason about but carry metadata overhead per character. Trap: proposing last-write-wins on the whole document — it silently discards a collaborator's paragraph.
Design a proximity service (Yelp "restaurants near me")
Central insight: unlike Uber (H12), the businesses barely move — so this is a read-heavy, write-rare geospatial problem, which means you can precompute aggressively. Store a geohash or S2 cell ID as an indexed column in a normal database, query the user's cell plus neighbours, then filter by exact distance and rank by rating/relevance. Cache per-cell result sets, since thousands of users in the same cell get the same candidate list. Trap: reaching for the Uber design. Recognising that a different write pattern permits a much simpler answer is the skill being tested.
Design a stock exchange / order matching engine
Central insight: the matching engine is deliberately single-threaded and in-memory — a limit order book as two priority queues (bids descending, asks ascending) with price-time priority. Determinism and microsecond latency matter more than horizontal scale, and one core can match millions of orders per second. Durability comes from an append-only input log replicated to hot standbys that replay the same deterministic sequence (event sourcing), not from a database write per trade. Shard by symbol, since orders for different symbols never interact. Trap: distributing the matching engine itself — that destroys the deterministic ordering that fairness and auditability depend on. This is the clearest example in the whole manual of a case where not distributing is the senior answer.
Part III — the LLD / machine-coding round
A machine-coding round gives you 60–120 minutes and a prompt like "design a parking lot" or "build Splitwise." You produce working, compiling Java plus a class design you can defend. It is graded far more on structure than on features.
What is actually scored
| Weight | Criterion | What it looks like |
|---|---|---|
| High | Working demo | A main that runs a realistic scenario end to end and prints sensible output. A beautiful design that does not run scores badly. |
| High | Extensibility | Adding a new vehicle type / payment method / discount rule touches one new class and zero existing ones (Open-Closed). |
| High | Separation of concerns | Entities, repositories, services, and strategies are distinct. No business logic inside a data class; no System.out.println inside a service. |
| Medium | Correct concurrency | If two threads can book the same seat, you handled it — and you can say why your approach is correct. |
| Medium | Naming and readability | ParkingSpotAssignmentStrategy, not Helper2. |
| Low | Feature completeness | Deliberately: a clean core with two features beats a sprawling mess with eight. |
The 90-minute plan
- 0–8 min — clarify and freeze scope. Write the requirements as a comment block at the top of the file. Explicitly list what is out of scope. This is your contract, and it protects you when time runs short.
- 8–20 min — identify entities and relationships. Nouns become classes, verbs become methods. Sketch the class diagram (on paper or as comments) and say which patterns you will use and why.
- 20–70 min — code, bottom-up. Enums and value objects → entities → repositories (in-memory) → services → strategies → the demo. Compile every 10 minutes.
- 70–85 min — the demo. A
maincovering the happy path and one edge case. - 85–90 min — narrate extensions. "To add a bike-only floor, I'd add one enum value and one strategy implementation; no existing class changes."
(1) In-memory repositories, always — use ConcurrentHashMap behind a repository interface. Nobody wants you to configure a database, and the interface shows you know where the seam belongs. (2) Never stop to make it perfect. A compiling, running system at minute 70 that you then improve beats an elegant half-system at minute 90. Get to "it runs" as fast as possible, then refactor.
Over-engineering. Candidates apply nine patterns to a problem that needs two, and run out of time. Use a pattern when it removes an if/else chain that will grow, or isolates something that genuinely varies. If you cannot name what varies, you do not need the pattern.
OOP foundations in Java
The four pillars, stated usefully
- Encapsulation — hide state behind behaviour. The test: can an external caller put the object into an invalid state? If a
BankAccountexposessetBalance(), it can. Exposedeposit()andwithdraw()instead, and validate inside. - Abstraction — expose what, hide how. A
PaymentGatewayinterface withcharge(amount)lets you swap Stripe for Razorpay without a caller noticing. - Inheritance — an "is-a" relationship for sharing behaviour. Powerful and over-used; see composition below.
- Polymorphism — one interface, many implementations, resolved at runtime. This is what lets you delete
if (type == CAR) ... else if (type == BIKE).
Composition over inheritance — the single most useful principle in LLD
"Use inheritance for is-a relationships where the subtype is genuinely substitutable, and composition for has-a or can-do relationships. In practice most 'is-a' relationships turn out to be 'behaves-like-a', which means composition." Then cite the classic counterexample: Square extends Rectangle looks correct mathematically and violates the Liskov substitution principle, because setting a rectangle's width independently of its height is meaningful for the parent and impossible for the child.
Interface vs abstract class
| Interface | Abstract class | |
|---|---|---|
| Multiple | A class can implement many | Only one may be extended |
| State | No instance fields (only static final constants) | Can hold fields and constructors |
| Methods | abstract, default, static, private | Any, including protected helpers |
| Means | A capability or contract: can-do | A partial implementation of a family: is-a |
| Use for | Strategies, listeners, repositories — anything you will have several unrelated implementations of | Shared state and a template method skeleton across closely-related subclasses |
Modern Java tip worth mentioning: sealed interfaces plus records give you a closed set of implementations the compiler can exhaustively check in a switch — very useful for modelling a fixed set of events, commands, or states.
Equals / hashCode — the detail that gets noticed
SOLID, with before and after
S — Single Responsibility
A class should have one reason to change. The practical test: can you describe what it does without saying "and"?
O — Open/Closed
Open for extension, closed for modification. In practice: a growing if/else or switch on a type is the smell, and polymorphism is the cure.
L — Liskov Substitution
A subtype must be usable anywhere its supertype is, without the caller knowing. Violations look like: throwing UnsupportedOperationException in an override, strengthening preconditions, or weakening postconditions.
I — Interface Segregation
No client should be forced to depend on methods it does not use. Many small role interfaces beat one fat one.
D — Dependency Inversion
High-level modules should depend on abstractions, not on concrete low-level classes. This is what makes code testable.
Do not recite the acronym. Narrate the principle at the moment you apply it: "I'm putting the fee calculation behind a strategy interface so adding a new fee type doesn't modify this class — that's the open-closed part." One sentence at the right moment is worth more than a definition.
Practice — Name a case where following SOLID would be wrong.judgement
When the variation you are abstracting over does not exist. Creating a PaymentMethod interface with exactly one implementation, forever, adds indirection and a file to navigate for no benefit — that is speculative generality, and it makes the code harder to read, not easier. The same applies to splitting a 40-line class into six classes because "single responsibility": if all six always change together, they had one responsibility all along.
The honest formulation: SOLID is a response to observed change, not a prophecy about it. Apply it where you can name the axis of variation ("we will certainly add more payment methods"), and inline it where you cannot. Interviewers respect a candidate who says "I'd keep this concrete until a second implementation appears" far more than one who abstracts everything — the second is the most common cause of unreadable LLD submissions.
Class diagrams: the five relationships
Aggregation vs composition is the one interviewers actually ask about, and the test is lifetime ownership. In a parking-lot design: ParkingLot composes Floor composes ParkingSpot (destroying the lot destroys the spots — they have no independent meaning), but ParkingSpot only associates with Vehicle (the car exists before and after it parks). Getting this right in your diagram signals real modelling ability.
For a 90-minute round you do not need a UML tool. A boxes-and-lines sketch with multiplicities, or even a comment block listing classes with their fields and relationships, is enough — what matters is that you can explain the relationships and defend the ownership choices.
Creational patterns
These control how objects come into existence, so that callers do not depend on concrete classes or on complex construction logic.
Singleton — exactly one instance
Singletons are global mutable state: they make unit testing hard (you cannot substitute a fake), they hide dependencies (a class that calls Db.getInstance() does not declare that it needs a database), and they create contention in concurrent code. In an LLD round, prefer one instance created at the composition root and injected — you get "exactly one" without any of the costs. Mention Singleton, then say you would use dependency injection instead; that answer is strictly better than reciting the pattern.
Factory Method and Abstract Factory
Builder — the one you will actually use
Prototype completes the family: when constructing an object is expensive (a heavy parsed config, a pre-populated game board), clone() an existing instance instead. In Java, prefer a copy constructor or a toBuilder() method over Cloneable, which has a famously broken contract — and say so if asked.
Structural patterns
These compose objects into larger structures without making the composition rigid.
| Pattern | Intent | Tell-tale sign you need it |
|---|---|---|
| Adapter | Make an incompatible interface fit | Integrating a third-party SDK whose API does not match yours |
| Decorator | Add behaviour to one object at runtime, layered | Optional, combinable features: pricing add-ons, stream compression + encryption |
| Facade | One simple entry point over a complex subsystem | A checkout flow touching inventory + pricing + payment + shipping |
| Composite | Treat a tree of objects and leaves uniformly | File systems, org charts, UI hierarchies, nested discounts |
| Proxy | Stand in for another object to control access | Lazy loading, caching, access control, remote calls, rate limiting |
| Flyweight | Share immutable intrinsic state across many objects | Millions of similar objects: characters in an editor, tiles in a map, chess pieces |
| Bridge | Split abstraction from implementation so both vary | Shapes × renderers, notifications × transports |
All three wrap another object; the difference is intent. Adapter changes the interface (the caller could not use the wrapped object before). Decorator keeps the interface and adds behaviour (and is designed to stack). Proxy keeps the interface and controls access to it — lazily, remotely, or conditionally — without adding new behaviour the client asked for. If someone asks "aren't they the same?", that one sentence is the answer.
Behavioural patterns
| Pattern | Intent | Appears in |
|---|---|---|
| Strategy | Interchangeable algorithms selected at runtime | Pricing, parking-spot allocation, matching, eviction policies |
| Observer | One-to-many notification of state change | Notifications, event buses, UI updates, cache invalidation |
| State | Behaviour changes with internal state; each state is a class | Vending machines, order/trip lifecycles, elevators, TCP |
| Command | Encapsulate a request as an object | Undo/redo, job queues, remote controls, transaction logs |
| Chain of Responsibility | Pass a request along handlers until one handles it | Middleware, logging levels, approval workflows, ATM cash dispensing |
| Template Method | Fixed algorithm skeleton, overridable steps | Report generation, data import pipelines, test setup/teardown |
| Iterator | Traverse a collection without exposing its structure | Custom collections, paginated API clients |
| Mediator | Objects communicate through a hub, not with each other | Chat rooms, air-traffic control, complex form validation |
| Memento | Capture and restore state without violating encapsulation | Undo, checkpoints, game saves |
| Visitor | Add operations to a class hierarchy without modifying it | AST processing, compilers, document exporters |
(1) Catch exceptions per listener, as above — otherwise a failing analytics listener rolls back a paid order. (2) Say when you would make notification asynchronous: if listeners do IO, publish to a queue instead of calling them inline, so the order does not wait on an email server. That is the exact moment Observer (in-process) becomes pub/sub (distributed) — connecting the LLD pattern to the HLD pattern is a memorable answer.
Why this beats a switch: every state class holds all the behaviour for that state, including which actions are illegal, so adding a MaintenanceState cannot break the others. And the return type makes transitions explicit — you can read the machine's graph by scanning the return statements.
A strong parking-lot solution uses Strategy (spot allocation) and maybe Factory (vehicle creation) — that is it. A strong Splitwise uses Strategy (split types). A strong vending machine uses State. Each pattern must be justified by a variation that actually exists in the requirements. Interviewers deduct for pattern-stuffing far more often than they deduct for pattern-absence.
Practice — Strategy and State look identical in code. What's the difference?classic
Structurally they are near-identical: an interface, several implementations, and a context that delegates. The differences are who chooses and whether the implementations know about each other.
Strategy: the client picks the algorithm ("use the nearest-spot policy"), the strategies are unaware of one another, and the choice usually does not change during an operation. The variation is about how to do one thing.
State: the object itself transitions, states typically decide and return the next state, and the set of legal operations changes with the state. The variation is about what the object can currently do.
A crisp test: if the implementations return or set the next implementation, it is State. If they just compute a result, it is Strategy.
Java concurrency for LLD rounds
Almost every machine-coding prompt has a hidden concurrency requirement — two cars arriving at the same spot, two users booking the same seat, two threads incrementing the same counter. You are not expected to build a thread pool from scratch; you are expected to pick the right primitive and justify it.
The primitives, ranked by how often you should reach for them
| Tool | Use for | Notes |
|---|---|---|
ConcurrentHashMap | Any shared map. Your default repository. | computeIfAbsent, merge, and putIfAbsent are atomic — use them instead of check-then-act. |
AtomicInteger / AtomicLong / AtomicReference | Counters, IDs, lock-free state swaps | compareAndSet is the building block for optimistic updates. |
synchronized | Short critical sections guarding an object's own fields | Simplest correct thing. Reentrant. Cannot be interrupted or time out. |
ReentrantLock | When you need tryLock(timeout), fairness, or multiple condition variables | Must be released in a finally block, always. |
ReadWriteLock / StampedLock | Read-heavy shared structures | Many readers, one writer. Only worth it when reads dominate heavily. |
BlockingQueue | Producer–consumer handoff | LinkedBlockingQueue, ArrayBlockingQueue. Bounded queues give you backpressure for free. |
ExecutorService | Running tasks in the background | Never new Thread() in a service. Always shut down in a lifecycle method. |
CompletableFuture | Composing async steps, parallel fan-out | allOf, thenCombine, exceptionally. |
CountDownLatch / Semaphore | Wait for N tasks; limit concurrent access to N permits | Semaphore is a neat model for "N parking spots" or "N connections". |
(1) Unbounded queues — an unbounded LinkedBlockingQueue turns backpressure into an OutOfMemoryError. (2) Swallowing InterruptedException — catch it, restore the interrupt flag, and exit. (3) Non-volatile stop flags — a plain boolean may never be observed as changed by another thread. (4) Locking across IO — holding a lock while making a network call multiplies contention by the network's latency and is a classic outage cause.
"I'm making the state transition itself atomic with a compare-and-set on the spot, rather than locking the whole lot, so two cars going to different floors never contend. The only serialisation is between two cars competing for the same spot, which is exactly the contention that must exist."
Design a parking lot
Requirements (freeze these in the first five minutes)
- Multi-floor lot; spot types for motorcycle, car, and truck; a vehicle may occupy a compatible spot.
- Park a vehicle → issue a ticket. Unpark → compute a fee and free the spot.
- Pluggable spot-allocation policy and pluggable pricing.
- Thread-safe: many entry gates operate concurrently.
- Out of scope: persistence, reservations, valet, number-plate recognition, real payment gateways.
"An electric-charging spot type is one enum value plus one line in compatibleWith. Weekend pricing is one new PricingStrategy. Reservations would add a RESERVED state to the spot's atomic reference rather than a boolean, so the CAS still gives me exactly one winner. Persistence would go behind a SpotRepository interface — the service already never touches a collection directly except through the floors it was constructed with."
Design an elevator system
Requirements
- N elevators serving F floors. Two request sources: an external hall call (floor + up/down) and an internal car call (destination floor from inside).
- A controller assigns hall calls to the best car; each car serves its own stop list efficiently (do not reverse direction unnecessarily).
- Pluggable dispatch policy; doors open/close; capacity limits; a maintenance mode.
- Out of scope: physics, real motors, emergency systems.
An elevator is not a FIFO queue. It is a two-set scheduler: while moving up, it serves the sorted set of stops above it; when that empties, it flips direction and serves the sorted set below. This is the elevator (SCAN) disk-scheduling algorithm. Saying "I'll keep two TreeSets — up-stops and down-stops — and drain them in order" is the whole design in one sentence.
Follow-up — Morning rush: everyone goes from the lobby upward. How does your design adapt?extension
This is exactly why dispatch is a strategy. Add an UpPeakStrategy that (a) parks idle cars at the lobby instead of leaving them where they last stopped, (b) prefers to fill one car before dispatching it rather than sending half-empty cars, and (c) biases the cost function so a car already at the lobby wins ties. Swap it in on a schedule or when the controller detects that most hall calls originate from floor 0 — the controller code does not change at all.
Mention two more real refinements: zoning (assign cars to floor bands so they do not all traverse the whole shaft) and destination dispatch (passengers enter their destination in the lobby, the system groups people going to nearby floors into the same car). Naming destination dispatch is a nice signal — it is what modern buildings actually use, and it changes the API: the hall call carries a destination, so HallCall gains a field and the strategy gets more information to work with.
Design BookMyShow (movie ticket booking)
The LLD counterpart to H13. Everything here hinges on one question: can two users book the same seat? Design so the answer is structurally no.
Requirements
- Cities → theatres → screens → shows (a movie on a screen at a time). Seats per screen with types and prices.
- Search shows; view seat availability; hold selected seats for 10 minutes; confirm after payment; auto-release expired holds.
- Thread-safe under concurrent booking. Out of scope: real payments, persistence, recommendations.
(1) "Locks are scoped per show, so two different shows never contend." (2) "compute on a ConcurrentHashMap is atomic, so the claim is a compare-and-set — there is no window between checking and taking." (3) "Holds are all-or-nothing: if I can't get all four seats together, I release the ones I got, because a party of four does not want two seats." That third point is a product observation inside a concurrency answer, and it is the kind of thing that gets remembered.
Follow-up — Payment succeeds but the hold expired one second earlier. What now?edge case
This is the real-world case candidates skip, and simply naming it earns credit. Options, in order of preference: (1) Prevent it — extend the hold when payment starts ("payment in progress" state that is not sweepable), so the window closes; the hold TTL should cover the payment provider's worst-case latency, not just user think-time. (2) Detect and recover — if confirmation fails after a successful charge, immediately issue a refund and tell the user honestly; log it for reconciliation exactly as in H14. (3) Opportunistic re-acquire — if the seats are still free (nobody took them in that second), re-hold and confirm; only refund if they were genuinely taken.
The framing that matters: "any two-phase operation across an expiry boundary has a race, so I need either a state that suppresses expiry during phase two, or a compensating action. I'd build both — prevention for the common case and reconciliation as the safety net."
Design Splitwise (expense sharing)
Requirements
- Users and groups. Add an expense paid by one or more users, split among participants.
- Split types: EQUAL, EXACT amounts, PERCENT. Extensible to more.
- Show per-user balances ("you owe Rohit ₹450"). Simplify debts within a group.
- Out of scope: multi-currency conversion, settlements via real payment rails, persistence.
₹100 split three ways is 3,333 + 3,333 + 3,334 paise, not 3,333.33 each. Every strategy above guarantees the shares sum exactly to the total, with the remainder distributed deterministically. Mentioning this unprompted — "I use integer paise and give the remainder to the first N participants so the books always balance" — is a small thing that reads as real financial-software experience.
Design an LRU cache (and LFU)
The requirement is O(1) get and put. That constraint alone dictates the design: a hash map for lookup, plus a doubly linked list for recency ordering, with the map storing node references so any node can be unlinked in constant time.
(1) Sentinels. The dummy head and tail mean no branch ever checks for null — that is why the code is short and why it is hard to get wrong. (2) Java already has this: new LinkedHashMap<>(cap, 0.75f, true) with removeEldestEntry overridden gives you an LRU in five lines, and Caffeine is what you would use in production. Say it after writing the manual version — it shows you know the ecosystem without dodging the exercise. (3) Thread safety: synchronized here serialises everything; a real concurrent cache uses striped locks or a lock-free ring buffer of access events replayed in batch, which is exactly what Caffeine does.
Two short ones: rate limiter and logger
Thread-safe token bucket
Two points to narrate: refill is computed from elapsed time rather than run by a background thread (one timer thread for a million keys would not scale), and the lock is per bucket so unrelated users never block each other. Then connect to the HLD version: "across many servers I'd move this state into Redis with the Lua script from H2, because otherwise N nodes each allow the full limit."
A logging framework (Chain of Responsibility + Strategy)
Writing to a file on the caller's thread means a slow disk slows your request handler. The production answer is an async appender: the logger puts the event on a bounded queue and a dedicated thread drains it to the sinks. Then state the trade-off explicitly — a bounded queue means that under extreme load you must either block the caller (safe, slow) or drop log events (fast, lossy), and for logs, dropping is usually correct. That is exactly the choice Log4j2's async appender exposes, and naming it lands well.
Rapid fire: seven more LLD prompts
For each: the entities, the pattern that carries the design, and the trap.
Tic-tac-toe / Connect Four
Entities: Board, Cell, Player, Move, Game, WinningStrategy. Pattern: Strategy for win-detection so a 3×3, an N×N, and Connect Four share one Game. Trap: checking all rows, columns, and diagonals after every move — O(N²). Keep running counts per row, column, and the two diagonals and check only the four affected by the last move: O(1). Interviewers ask for exactly this optimisation.
Snake and Ladder
Entities: Board (a map of jump start → end), Dice, Player, Game with a turn queue. Pattern: Strategy for dice (a real die vs a loaded die for testing — this makes the game deterministically testable, which is the point). Trap: modelling snakes and ladders as two separate classes; they are the same thing — a jump — differing only in direction.
Chess
Entities: Board, Cell, abstract Piece with subclasses, Move, Game, Player. Pattern: polymorphic canMove(board, from, to) per piece — the textbook use of inheritance. Trap: forgetting the moves that break the "piece decides" model: castling, en passant, and promotion involve two pieces or a state change, so they belong on Game/MoveValidator, not on Piece. Naming that boundary problem is what separates a good answer here.
In-memory file system
Entities: FileSystemNode with FileNode and DirectoryNode. Pattern: Composite (L6) — size(), delete(), and search() recurse uniformly. Trap: path resolution. Write a single resolve(String path) that splits on / and walks the tree, and reuse it everywhere; candidates who re-parse paths inside every method produce a mess.
Key-value store with TTL
Entities: Store wrapping a ConcurrentHashMap<String, Entry> where Entry holds value + expiresAtNanos. Design point: combine lazy expiry (check on read and treat expired as absent) with active expiry (a background sweeper sampling random keys, exactly like Redis). Lazy alone leaks memory for keys never read again; active alone is too expensive to be exhaustive. Trap: using a Timer per key — a million keys means a million timers.
Notification service
Entities: Notification, Channel, UserPreferences, NotificationService. Patterns: Observer (subscribers react to events) + Strategy (one sender per channel) + Factory (pick the sender). Trap: synchronous sending inside the business transaction — if the email server hangs, the order hangs. Queue it, as in H6.
Uber cab matching (LLD scope)
Entities: Rider, Driver, Trip, Location, TripState, MatchingStrategy. Patterns: State for the trip lifecycle (REQUESTED → ASSIGNED → STARTED → COMPLETED, with only legal transitions), Strategy for matching (nearest, highest-rated, surge-aware), Observer for trip-status updates to both parties. Trap: the double-assignment race — use an atomic compare-and-set on the driver's state exactly as in the parking lot (L9), not a mutable boolean.
LLD pitfalls and a pre-submit checklist
| Pitfall | What it looks like | Fix |
|---|---|---|
| God class | BookingManager with 900 lines doing search, pricing, payment, and email | Split by reason-to-change; a service should orchestrate, not compute |
| Anaemic model | Entities are pure getters/setters; all logic sits in services | Put invariants on the entity — spot.tryOccupy(), not spot.setOccupied(true) from outside |
| Primitive obsession | String userId, String showId, long amount everywhere | Small value types (Money, SeatId) — the compiler then catches argument-order bugs |
| Pattern stuffing | Nine patterns in a 200-line solution | Use one only when you can name the axis of variation |
| Ignoring concurrency | Check-then-act on shared state | Atomic map ops or CAS; scope locks to the contended entity |
Money as double | 0.1 + 0.2 != 0.3 | Integer minor units (paise/cents) or BigDecimal |
| No demo | Beautiful classes, nothing runs | Write main at minute 60, not minute 89 |
| Unbounded growth | Maps that only ever grow | Eviction, TTL, or an explicit bound — and say which |
Five-minute pre-submit checklist
- It compiles, and
mainruns a happy path plus one failure path. - Requirements and out-of-scope items are written at the top as comments.
- Every interface has at least two conceivable implementations (otherwise inline it).
- No business logic in entity setters; no
System.out.printlninside services. - Shared mutable state is either confined, atomic, or guarded — and I can say which.
- Money is integer minor units; times are
Instant, injected rather thanInstant.now()deep inside logic. - Exceptions are specific (
SeatsUnavailableException), not bareRuntimeException. - Collections returned from getters are copies or immutable views.
- I can name the two extensions I would make next and which class each touches.
"The core is a service that orchestrates, entities that protect their own invariants, and two strategies for the parts the requirements said would vary. If you want, I can show how adding <the extension they hinted at> fits — it's one new class and no changes to existing ones." Then stop talking. Ending with a demonstration of extensibility is the strongest possible finish.
Part IV — a mock HLD round, annotated
What follows is a compressed 45-minute round. The annotations in the margin notes are what the interviewer is thinking.
Good. He bounded an unbounded question in 30 seconds and asked for confirmation instead of assuming. Candidates who start drawing here almost always run out of time.
He ended the estimation with a conclusion — "the problem is the read path, and photos never touch the DB." That one sentence tells me he understands why he did the math. Most candidates recite numbers and then design as if they had not.
The "only materialise for active users" line is the kind of thing that separates candidates. It shows he did the memory arithmetic in his head and found the cheap fix instead of accepting a 4 TB Redis bill.
He volunteered the weaknesses before being asked, named the fallback path as something that must be tested, and framed ranking as a layer that does not disturb the storage design. Notice also what he never did: he never listed technologies for their own sake, and every number he computed changed a decision.
A mock machine-coding round, annotated
Two minutes spent writing the contract is the best-spent time in the whole round. It also means that if he runs out of time at minute 85, I can see what he intended, and I grade that.
He did not say "I'd add synchronized." He identified which resource is contended, chose the narrowest primitive, and described the failure the loser sees. The user-visible consequence is half the answer.
Naming the one weakness you did not have time to fix, precisely, is far stronger than hoping the interviewer misses it. It signals you know what good looks like even when the clock beat you — and it pre-empts the criticism they were about to write down.
Company flavours: what each type actually weights
| Company type | Format | What they weight | Prepare by |
|---|---|---|---|
| US big tech (Google, Meta, Amazon, Microsoft) | 45–60 min HLD whiteboard, usually at SDE-2+; SDE-1 may get a lighter version | Structured thinking, trade-off articulation, scale reasoning, communication. Amazon layers leadership principles onto it. | H1–H12 out loud, timed. Practise narrating while drawing. |
| Indian product companies (Flipkart, Swiggy, Zomato, PhonePe, Razorpay, Atlassian, Uber India, Navi) | Machine coding round is often decisive — 90–120 min, laptop, working code; then a separate HLD round | Compiling code, extensibility, SOLID, correct concurrency, clean naming | Part III. Type L9, L11, L12 from scratch, timed, without notes. |
| Startups / early-stage | Practical, often about their actual problem | Pragmatism, cost awareness, "what would you build this week" | Be ready to defend not distributing something. F19 and H13's "the contended state is tiny" framing. |
| Service companies & core-CS rounds | Rapid-fire conceptual Q&A | Definitions, DBMS/OS/CN fundamentals, patterns by name | X7 below, plus your OS/DBMS/CN material. |
| Quant / HFT / infra | Latency-obsessed, systems-level | Memory layout, lock-free structures, determinism | H18's matching engine; know why single-threaded can beat distributed. |
On-campus processes compress everything: an online assessment, then two or three rounds in a day, often with a mixed DSA + design round rather than a dedicated one. That means breadth beats depth — being able to give a competent 15-minute answer to any of H1–H12 is worth more than a brilliant 45-minute answer to one. Off-campus and experienced-hire loops go deeper and will push on one component until you reach the edge of your knowledge; there, the right move is to say "I don't know that layer well, here's how I'd find out" rather than improvising confidently.
Phrases that score, and red flags
- "Before I design anything — is this read-heavy or write-heavy, and are we optimising for availability or consistency?"
- "That's about 3,000 writes a second, which one Postgres node handles, so I won't shard yet. Here's what would make me shard."
- "I'd start with a monolith and extract this service when a second team owns it."
- "The trade-off is X versus Y; I'm choosing X because our requirement said Z."
- "This is eventually consistent, and the user-visible consequence is …"
- "The single point of failure here is …, and I'd fix it by …"
- "If traffic grew 10×, the first thing to break would be …"
- "I'd measure … before optimising this."
- "I don't know that in detail — here's how I'd reason about it / find out."
- "Let me check: is it more useful to go deeper here, or cover the read path?"
- Drawing Kafka, Redis, and Cassandra before asking a single question.
- "We'll use Cassandra because it scales." (Why? What access pattern?)
- "Microservices" as a default for a system with 100 users.
- Computing QPS to three decimal places and then ignoring it.
- Saying "eventually consistent" without saying what the user sees.
- Ignoring a hint, or saying "good point" and changing nothing.
- Silence for five minutes while you think — narrate instead.
- Claiming certainty about something you half-remember.
- Adding a cache to every box on the diagram.
- In LLD: nine patterns, no working
main.
Design interviews are graded on how you decide, not on what you decide. Two candidates can propose the same architecture and score three levels apart, because one derived it from stated requirements and one recited it. Every time you make a choice, attach a reason to it in the same breath — that habit alone is worth more than another five case studies.
One-page revision sheet
Study plans
Four weeks (interview is close)
| Week | Read | Do |
|---|---|---|
| 1 | F1–F10 | 10 estimation drills. Draw Fig. 0 from memory daily until it takes 60 seconds. |
| 2 | F11–F19, H1–H6 | One case study per day: 20 min solo on paper first, then diff against the text. |
| 3 | L1–L8 | Type L9 (parking lot) and L11 (BookMyShow) from scratch, timed at 75 min, no notes. |
| 4 | H7–H18, X1–X5 | Three mock interviews out loud with a friend or a recording. Re-do your two worst case studies. |
Eight weeks (building real depth)
Same as above at half pace, plus: implement one thing per week for real — a consistent hash ring with a rebalancing test, a token-bucket limiter backed by Redis, an LRU cache benchmarked against LinkedHashMap, a small Kafka consumer that is genuinely idempotent. Building beats reading by a wide margin, because the details you get wrong are exactly the details interviewers probe.
Record yourself answering a case study for 20 minutes, then listen. It is uncomfortable and it is the fastest feedback loop available: you will immediately hear the filler, the unexplained jumps, and the places where you stated a technology instead of a reason. Two recordings are worth ten passive re-reads.
Rapid-fire Q&A
Cover the answer, say yours out loud, then check. Aim for two sentences each.
Difference between latency and throughput?basics
Latency is time per operation; throughput is operations per unit time. They are not inverses — batching and pipelining raise throughput while increasing per-request latency, which is exactly the trade a queue makes.
Horizontal vs vertical scaling — when is vertical right?basics
Vertical is right when you are below the ceiling of a single machine and want to avoid distributed-systems complexity — which covers far more real systems than people admit. It stops being right when you need redundancy or exceed the largest instance available.
What does a load balancer do that a DNS round robin doesn't?networking
Health checking, session handling, TLS termination, and content-based routing. DNS has no idea a backend is dead and clients cache the answer past its usefulness.
Why is a CDN cheaper than scaling your origin?networking
Because the same bytes are served to many users from a location close to them, so you pay for one origin fetch instead of millions, and you avoid long-haul bandwidth. It also absorbs spikes and some DDoS for free.
When would you choose UDP over TCP?networking
When late data is worthless and retransmission hurts more than loss: real-time voice and video, game state, metrics. Also as QUIC's substrate, where reliability is rebuilt per-stream to avoid TCP head-of-line blocking.
SSE vs WebSocket?real-time
SSE is one-directional server→client over plain HTTP with automatic reconnection — simpler to operate. WebSocket is full duplex and is what you need when the client sends frequently too.
What breaks when you add read replicas?databases
Read-your-own-writes. A user posts, the read hits a lagging follower, and their own post is missing. Fix with leader reads for a short window after a write, or replica pinning.
B-tree vs LSM tree — pick one and justify.databases
B-tree for read-heavy workloads with in-place updates and range scans; LSM for write-heavy append-mostly workloads, accepting read amplification and compaction cost. The workload's read:write ratio decides it.
Why does a random UUID primary key hurt?databases
Random inserts scatter across the clustered index, causing page splits and poor cache locality, so write throughput drops and the index grows. Use a time-ordered ID (Snowflake, UUIDv7, ULID) instead.
Explain the leftmost prefix rule.databases
A composite index on (a,b,c) can serve queries filtering on a, or a+b, or a+b+c — but not on b alone, because the index is sorted by a first. Column order in a composite index is a design decision, not a formality.
What is write skew and which isolation level prevents it?transactions
Two transactions each read a set of rows, make disjoint writes, and together break an invariant across those rows — no single row was written twice, so snapshot isolation misses it. Serializable prevents it; so does materialising the conflict into a row both must write.
Optimistic or pessimistic locking for a hot inventory row?transactions
Pessimistic. Optimistic locking degrades into a retry storm precisely when contention is high, which is the case you are trying to survive.
How do you avoid deadlock in a money transfer?transactions
Acquire the two account locks in a deterministic order — ascending account ID — so no cycle can form. Better still, use a single conditional statement or an append-only ledger so there is nothing to lock.
Sharding key for chat messages, and why?sharding
chatId, so an entire conversation lives on one shard and "load the last 50 messages" is a single-partition read. Very large group chats get sub-partitioned by time bucket to bound partition size.
What is a hot partition and give three fixes.sharding
One shard receiving disproportionate traffic because a key is popular. Fixes: cache the key, salt it across N sub-keys and scatter-gather on read, or special-case it out of the normal path entirely.
Why virtual nodes in consistent hashing?hashing
With few physical nodes the ring is uneven, and a departing node dumps all its load on one neighbour. Many virtual positions per node even out ownership and spread a departure across the whole cluster.
State CAP correctly.distributed
Partitions are not optional, so the real statement is: during a partition you must choose availability or consistency. PACELC adds the part that matters daily — even without a partition, you trade latency against consistency.
What does W + R > N give you?distributed
Guaranteed overlap between the write set and the read set, so a read sees the latest acknowledged write. N=3, W=2, R=2 is the standard configuration.
Why is last-write-wins dangerous?distributed
It depends on synchronised wall clocks; a few milliseconds of skew silently discards a genuinely later write. Use version vectors or logical clocks to detect concurrency instead of guessing at it.
Is exactly-once delivery possible?messaging
Not as an end-to-end network guarantee. You get at-least-once delivery plus idempotent processing, which is effectively once — and that is what to say.
What problem does the transactional outbox solve?messaging
Writing to a database and publishing an event atomically. You insert the event into an outbox table in the same local transaction, then a relay (usually CDC) publishes it, so the event exists if and only if the data was committed.
Kafka: what determines ordering?messaging
The partition key. Messages with the same key land on the same partition and are strictly ordered relative to each other; there is no global order across partitions.
Cache-aside vs write-through — which and why?caching
Cache-aside by default: simple, and only requested data occupies memory. Write-through when staleness is unacceptable and you can afford both write latencies.
Why delete rather than update a cache entry on write?caching
Concurrent writers can interleave and leave the older value permanently in the cache. Deleting forces the next read to re-populate from the source of truth.
What is a cache stampede and how do you stop it?caching
A popular key expires and thousands of concurrent misses hit the database at once. Stop it with single-flight request coalescing, TTL jitter, or probabilistic early refresh.
Why add jitter to retries?reliability
Because synchronised clients retry in lockstep and re-create the exact spike that caused the failure. Full jitter decorrelates them and lets the dependency recover.
What does a circuit breaker actually buy you?reliability
It converts a slow cascading failure into a fast contained one, freeing threads that would otherwise pile up on a dead dependency. It also gives the dependency room to recover instead of being hammered.
Serial dependencies and availability — do the math.reliability
Five services at 99.9% each give 99.5% overall, worse than any single component. Availability multiplies along a serial path, which is an argument for fewer hops and for graceful degradation.
2PC or Saga across microservices?distributed
Saga. Two-phase commit blocks all participants when the coordinator fails and holds locks across services; a saga is a sequence of local transactions with compensations, which keeps each step independently available.
Why are distributed locks unsafe, and what fixes them?distributed
A process can pause (GC) past the lock's expiry and resume believing it still holds it. Fencing tokens fix it: storage rejects any write carrying a token older than the highest it has seen.
Why odd-numbered Raft clusters?consensus
2f+1 nodes tolerate f failures, so 3 and 4 both tolerate only one — the fourth node adds cost and latency for zero extra fault tolerance.
Average latency vs p99 — why does it matter?observability
Averages hide the tail, and the tail is what users notice. Worse, a page making ten parallel calls turns a 1% slow-call rate into roughly 10% slow page loads.
What is an error budget for?observability
It turns reliability into a quantity you can spend: budget remaining means ship features, budget exhausted means stop and fix. It converts an argument into arithmetic.
Sessions or JWTs?security
Short-lived JWT access tokens for stateless verification, plus revocable server-side refresh tokens. Pure JWTs cannot be revoked before expiry, which is usually unacceptable for logout.
Token bucket vs leaky bucket?security
Token bucket allows bursts up to its capacity, which is friendlier for real API clients. Leaky bucket enforces a perfectly smooth output rate at the cost of queueing latency.
Composition or inheritance, and why?lld
Composition, by default. Inheritance couples you to a parent's implementation and explodes combinatorially when capabilities are independent; composition lets behaviour change at runtime.
Give a Liskov violation in one line.lld
Penguin extends Bird overriding fly() to throw. Model the capability (Flyable) rather than the taxonomy.
Strategy vs State — the one-line distinction?lld
If the implementations decide and set the next implementation, it is State. If they only compute a result chosen by the client, it is Strategy.
Why is double-checked locking broken without volatile?java
Without volatile, instruction reordering allows another thread to observe a non-null reference to a partially constructed object. volatile establishes the happens-before edge that prevents it.
What is wrong with an unbounded LinkedBlockingQueue?java
It removes backpressure: producers never block, so a slow consumer turns into unbounded memory growth and an OutOfMemoryError instead of a visible slowdown.
Why must equal objects have equal hash codes?java
Hash-based collections locate a key by bucket before comparing with equals. Break the contract and a HashMap silently fails to find entries that are logically present.
Why never use double for money?java
Binary floating point cannot represent decimal fractions exactly, so sums drift and comparisons fail. Use integer minor units or BigDecimal.
How do you make an at-least-once consumer safe?applied
Give every event a unique ID and claim it with an insert that has a unique constraint, in the same transaction as the side effect. Redelivery then becomes a no-op.
Design constraint: 100k users, 20k seats. What does that tell you?applied
That the contended state is tiny, so keep it on one node with real transactions and spend the scaling effort on the read path and an admission queue. Distributing the consistent part would solve the wrong problem.
How would you keep a search index in sync with a database?applied
Stream changes from the database's log (CDC) or an outbox table into the index, and accept sub-second staleness. Dual-writing from application code will drift, and you will need a reconciliation job either way.
Your p99 is fine but users complain. What's your first hypothesis?applied
You are measuring the wrong thing — server-side latency excludes DNS, TLS, queueing, client render, and failed requests. Measure end-to-end from the client and split successful from failed requests.
Question bank — practise these out loud
- Design a system that stores 1 PB of logs and answers "show me errors from service X in the last hour."
- Your single Postgres is at 90% CPU. Walk through your options in order.
- Design multi-tenant storage where one tenant must never affect another's performance.
- Migrate a 4 TB table to a new schema with zero downtime.
- Design a seat/inventory system that must never oversell, at 50k requests/sec.
- Two datacenters, both accepting writes. How do you resolve conflicts?
- A user updates their profile; three services cache it. How do they learn?
- Design an idempotent "transfer money" API.
- Design live comments on a stream with 1M concurrent viewers.
- Design presence for a 500-person team chat.
- Design a collaborative whiteboard.
- Design push notification delivery guaranteeing at-most-one per event.
- Parking lot, elevator, vending machine, ATM.
- Splitwise, BookMyShow, food delivery order flow.
- LRU/LFU cache, rate limiter, logging framework, in-memory file system.
- Chess, snake and ladder, tic-tac-toe with an N×N win check.
- A thread-safe bounded blocking queue, from scratch.
Set a 25-minute timer. Talk out loud the whole time — to a wall if necessary. Draw. When the timer ends, write down the one thing you skipped and the one decision you could not justify. Those two items are your actual study list; everything else is revision you have already done.