The non-AI design round — distributed systems, graded at staff. The six-step drive script, the estimation numbers you carry in, a primitive toolbox (sharding, quorums, consensus, back-pressure, idempotency), and six worked designs chosen because a datacenter AI-infra org actually asks them: rate limiter, rack telemetry, job scheduler, fleet control plane, artifact distribution, log and trace storage.
There are two different design rounds in this loop and candidates routinely prepare for only one of them.
For a datacenter org this round is not a formality. Qualcomm is standing up rack-scale infrastructure — an Infrastructure Management Suite for fleets of AI200 racks, on-prem appliances, hyperscaler deployments — and the software problems are ordinary, enormous distributed-systems problems. The six designs in §5 are all chosen on that basis: each is a thing this org has to have built.
The bar, in one sentence: at IC4 you produce a working architecture; at IC5 you produce a working architecture plus the estimate that sized it, the failure mode that will actually page you, and the blast-radius control that keeps one mistake from becoming an outage.
The one habit that carries this round: every design decision gets a number or a named alternative attached to it. "We'll use a queue" is IC3. "We'll use a queue with a bounded depth of 10,000 and shed above it, because at 2,000 jobs/s a deeper queue just converts a capacity problem into a latency problem for everyone" is IC5.
The AI design round Back to Round 1
Six moves, 45 minutes. Unlike the AI round, this one has a canonical order and interviewers expect it.
| Min | Move | Detail |
|---|---|---|
| 0–5 | Requirements | Functional (what it does) and non-functional (scale, latency, consistency, availability). Write both lists. Explicitly ask which non-functional property is the priority — it decides every later fork. |
| 5–9 | Estimates | QPS, storage/day, bandwidth. Out loud, with stated assumptions. §3 has the numbers. |
| 9–13 | API | Three to five endpoints with their parameters. This pins the data model and surfaces ambiguity cheaply. |
| 13–17 | Data model | Entities, keys, access patterns → then pick the store. Never the other way round. |
| 17–27 | High-level design | Boxes and arrows, then trace one write and one read end to end. |
| 27–38 | Deep dive ×2 | You pick. Say so out loud: "the two things that decide whether this works are X and Y." |
| 38–45 | Bottlenecks, failure, scale | What breaks at 10×, what pages you, what the blast-radius control is. |
Three phrases that read as staff.
Latency ladder (order of magnitude is what matters, not precision):
| Operation | Time | Rule of thumb |
|---|---|---|
| L1 cache | ~1 ns | |
| L3 cache | ~30 ns | |
| Main memory | ~100 ns | ~100× slower than L1 |
| NVMe random read | ~50–100 µs | ~1,000× slower than RAM |
| Same-DC network round trip | ~0.5 ms | ~5× slower than NVMe |
| Disk seek (spinning) | ~10 ms | |
| Cross-region round trip | ~50–150 ms | ~100,000× slower than RAM |
Throughput and conversion
The estimation move. Compute three numbers and say which one is the binding constraint: QPS, storage per day, bandwidth. One of them is always the problem, and naming which is the staff signal.
You should be able to deploy each of these in one sentence, with its cost.
Sharding and consistent hashing. Modulo hashing (hash(k) % N) remaps ~all keys when N changes — catastrophic for a cache. Consistent hashing places nodes on a ring and moves only ~1/N of keys on a membership change. What it still gets wrong: uneven load, which is why you use virtual nodes (100–200 tokens per physical node) to smooth the distribution, and why a single hot key defeats it entirely regardless.
Replication and quorums. With N replicas, W write acks, R read acks: W + R > N guarantees a read sees the latest write. W=N, R=1 favors reads; W=1, R=N favors writes; W=R=⌈(N+1)/2⌉ balances. W=1, R=1 is fast and eventually consistent. The interview move is picking different values for different stores in the same design and justifying each.
Consensus. Raft/Paxos gives you a single agreed-upon value across a majority. Use it for the things that must be unique and correct — leader election, cluster membership, config version. Do not put it in a high-throughput data path; it costs a majority round trip per decision.
CAP, honestly. Partitions happen, so the real choice is availability vs consistency during a partition. PACELC adds the part people forget: even when there's no partition, you're trading latency against consistency. That's the more useful framing for a datacenter design.
Caching. Cache-aside (app manages it; the default), write-through (consistent, slower writes), write-back (fast, can lose data). The two hard parts are invalidation and the thundering herd when a hot key expires — fix with request coalescing (single-flight) plus jittered TTLs.
Queues and back-pressure. A queue converts a throughput problem into a latency problem; it does not create capacity. A queue that is persistently non-empty is a mis-sized system, not a working one. Bound the depth and shed above it.
Delivery semantics. At-most-once loses messages. At-least-once duplicates them. Exactly-once end-to-end does not exist across a network — what exists is at-least-once delivery plus idempotent consumers, which is observationally equivalent and is the only correct answer to that question.
Idempotency. A client-supplied key plus a dedupe table with a TTL. Every mutating API in every design below needs one; naming it unprompted is a reliable staff marker.
Storage engines. LSM trees (RocksDB, Cassandra) — fast writes, compaction cost, reads may touch multiple levels; use bloom filters to skip them. B-trees (Postgres, MySQL) — balanced reads and writes, better for range and transactional work. Write-heavy fleet telemetry is LSM territory; a config store is B-tree territory.
Failure isolation. Timeouts (always), retries with jitter (never without — synchronized retries are a self-inflicted DDoS), circuit breakers, bulkheads, and graceful degradation defined in advance as a ladder, not improvised during the incident.
"Forty gateway nodes, one global limit per API key. Sub-millisecond overhead."
The most-asked classic, and it connects directly to B2 in the coding page — that's the single-node algorithm, this is the distributed system around it.
Requirements. Per-key limits, multiple tiers, global (not per-node) enforcement, <1 ms added latency, and — ask this — fail open or fail closed?
Algorithm choice, in the order you should present it:
| Algorithm | Memory | Burst behavior | Verdict |
|---|---|---|---|
| Fixed window | O(1) | 2× the limit across a boundary | Reject, and say why |
| Sliding window log | O(requests) | Exact | Correct but too expensive |
| Sliding window counter | O(1) | Approximate, smooth | Good compromise |
| Token bucket | O(1) | Controlled burst up to capacity | Pick this |
| Leaky bucket | O(1) | Smooths output, queues input | For traffic shaping, not admission |
Distribution — three options, and you should name all three before choosing.
rate/N per node. Zero latency, no dependency — but unfair under skewed routing (a node with no traffic hoards its share) and it under-admits.# Option 1's critical section — must be ONE atomic script, not three round trips
def check(key, cost, now):
b = load(key) # {tokens, last}
b.tokens = min(b.cap, b.tokens + (now - b.last) * b.rate)
b.last = now
if b.tokens >= cost:
b.tokens -= cost
store(b); return ALLOW
store(b)
return DENY(retry_after=(cost - b.tokens) / b.rate)Deep dives.
K sub-buckets at rate/K, chosen by hashing the request ID.Response contract. 429 with Retry-After, plus X-RateLimit-Limit / -Remaining / -Reset. Returning a bare 429 with no guidance guarantees clients hammer you.
Pushback: "Why not do it at the load balancer?" — Because limits are per-tenant, not per-IP, and the LB can't see the authenticated identity or the token cost. Edge rate limiting by IP is a separate, additional layer for DDoS, and you want both.
"1,250 racks, 64 accelerators each, 200 metrics per device, once a second."
Estimate first — and the number is the point of the question.
devices = 1,250 racks × 64 = 80,000
points/sec = 80,000 × 200 × 1 Hz = 16,000,000 points/s
raw @ 16 B = 256 MB/s = 22 TB/day <-- unaffordable
compressed ≈ 1.5–2 B/point (delta-of-delta timestamps + XOR values)
≈ 30 MB/s ≈ 2.6 TB/day <-- affordable, and that's the design
active series = 80,000 × 200 = 16,000,000Then say the thing that decides the design: cardinality kills you before volume does. Sixteen million active series is already large; add one per-request label and you multiply it by the number of distinct values. Budget series count explicitly and reject high-cardinality labels at ingest.
Architecture.
device agent ──push/scrape──> regional collector (shards by device hash)
│
buffer (Kafka: absorbs bursts, replays on failure)
│
┌──────────────┼──────────────┐
stream aggregator TSDB writer alert evaluator
(rollups, 1s→1m) (LSM, sharded) (SEPARATE path)
│
query + dashboardsThe five decisions worth defending.
Failure mode to walk: a deploy adds a label containing a request ID. Series count goes from 16M to 400M in ten minutes, the TSDB's index blows out, and ingestion stalls for everyone. Detection: alarm on series-creation rate, not just series count — the derivative gives you ten minutes of warning. Mitigation: per-tenant series quota enforced at ingest, and drop-with-a-loud-metric rather than accept-and-die.
"Submit, schedule, retry, prioritize, cancel. Thousands of workers."
Requirements to pin. Delayed/scheduled jobs or just a queue? (Both.) Priorities? (Yes.) At-least-once or at-most-once? (At-least-once.) Must a job survive a worker dying mid-execution? (Yes — this is the core of the design.)
Why not just a message queue? Because you need scheduled execution, priority, status query, and cancellation — a queue gives you none of those. You need a job store plus a claim protocol; the queue is an implementation detail underneath.
The claim protocol — this is the deep dive.
# Workers LEASE jobs. A lease is a timed claim, not a delivery.
def claim(worker_id, now):
job = atomically_select_and_update(
where = (status == READY and run_at <= now)
or (status == RUNNING and lease_expires < now), # reclaim the dead
order = (priority DESC, run_at ASC),
set = {status: RUNNING, owner: worker_id,
lease_expires: now + LEASE_TTL, attempts: attempts + 1},
limit = 1)
return job
def heartbeat(job_id, worker_id, now):
# extend ONLY if we still own it — a reclaimed job must not be re-extended
atomically_update(where = (id == job_id and owner == worker_id),
set = {lease_expires: now + LEASE_TTL})
def complete(job_id, worker_id, result):
atomically_update(where = (id == job_id and owner == worker_id),
set = {status: DONE, result: result})Say all four of these.
attempts bounds retries. Past MAX_ATTEMPTS, move to a dead-letter queue with the error — never retry forever, and never silently drop.owner check on every update is what prevents a slow worker from completing a job that was already reclaimed and re-run. Without it you get two writers for one job.(job_id, attempt) or, better, from the job's business key.Priority and fairness. Straight priority ordering starves low-priority work forever. Use weighted fair queues per tenant with an aging term that raises a job's effective priority with wait time — the same deficit-round-robin idea as the inference scheduler.
Back-pressure. Bound the ready queue. Above the bound, reject submissions with a retryable error rather than accepting work you cannot do. Watch job lag (age of the oldest ready job) as the primary health metric — rising lag means the worker pool is undersized, and it's the number that should drive autoscaling.
Pushback: "Just use SQS/Kafka." — Reasonable for the transport, but neither gives you scheduled execution, priority reordering, per-job status, or cancellation. I'd use the queue for dispatch and keep the job store authoritative.
"Push configuration and software to 1,250 racks without ever taking the fleet down."
The reframe that gets you the level: this is not a push system, it is a reconciliation loop. The control plane publishes desired state; agents pull it, converge toward it, and report observed state. Push-based control planes fail badly — a rack that was offline during the push silently stays stale forever, and nobody finds out until it matters.
control plane rack agent (every 30s)
desired_state (versioned, signed) ──> fetch desired
observed_state <────────────────── compare to local
health rollup converge, or roll back on health-check fail
report observed + healthFour properties, four sentences.
Consistency choice. The desired-state store needs strong consistency (linearizable, Raft-backed) — two racks reading different "current" versions is exactly the split-brain you're trying to avoid. The observed-state store is eventually consistent and high-volume; it's telemetry, and it belongs in D2's pipeline.
The failure to walk — the one that actually happens. A config change breaks the agent's ability to fetch the next config. Now you have bricked every rack the rollout reached, and the fix cannot be delivered by the mechanism that broke. Defenses, in order: the agent validates and health-checks before committing; last-known-good rollback is automatic; the fetch path is deliberately excluded from configurable surface area; and there is an out-of-band recovery channel. Any rollout system without a rollback path that's independent of the rollout path is one bad push from a site visit.
"A new 50 GB model binary has to reach 80,000 devices."
Start with the number that reframes the problem.
naive: 50 GB × 80,000 devices = 4,000,000 GB = 4 PB per release
at 10 Gbps from a central store: 4 PB / 1.25 GB/s ≈ 37 daysSay that out loud and the design writes itself.
The fix, in three layers.
1,250 × 50 GB = 62 TB instead of 4 PB — a 64× reduction — and the final hop is local bandwidth you already own. Add peer-to-peer between rack caches (Dragonfly/BitTorrent-style) if the rack count grows.Two more things to name.
Pushback: "Just use a CDN." — For the WAN hop, yes, and I would. It doesn't solve the last hop (rack cache → 64 devices over the LAN) or the rollback-staging problem, and for on-prem air-gapped appliances there is no CDN. The rack-cache tier is doing the work either way.
"3 TB a day of logs and traces. Almost nobody reads them — until something breaks, and then five people need them at once."
The economics are the design. Write-heavy, read-rare, bursty-read. Every decision follows from that asymmetry.
Architecture. Agents → Kafka → two divergent paths: an index over a small set of fields (service, level, trace ID, timestamp, a few tags) and the payload written as compressed blocks to object storage. Queries hit the index to find block offsets, then fetch only those blocks.
Say this: do not index the payload. Full-text indexing 3 TB/day costs more than the storage and buys you a search you use twice a month. Index the fields you filter on; grep the blocks you fetch.
Tiering with real numbers. Hot on SSD, 7 days (~21 TB, the incident window). Warm on object storage, 30 days. Cold in archival, 1 year, with a restore SLA measured in hours and stated up front so nobody expects otherwise during an incident.
Sampling — the staff answer. Head sampling (decide at trace start, keep 1%) is cheap and throws away exactly the traces you needed, because errors and slow requests are rare by definition. Tail sampling buffers a trace until it completes, then keeps it if it errored, was slow, or is a random baseline sample. It costs a buffering tier and it is worth it: you keep 100% of errors, 100% of the p99, and 1% of the boring ones. Retention drops ~50× with the interesting data intact.
Failure mode to walk: a service starts logging a stack trace per request at ERROR. Log volume goes 20×, Kafka lags, and every other service's logs are delayed — so the one team that needs logs right now can't see them. Fix: per-service ingest quotas with drop-and-count, plus rate-limited log deduplication at the agent (log the first N occurrences per minute, then a count). Quotas at ingest are the load-shedding pattern from D1, applied to a different resource.
| They say | You say |
|---|---|
| "Why not put everything in Postgres?" | "For the config store, I would — it's small, transactional, and needs strong consistency. For 16M metric points/s it's the wrong engine: that's an append-only, write-heavy, time-ordered workload, which is what LSM-based TSDBs exist for. I'd rather run two right stores than one wrong one." |
| "Exactly-once delivery." | "Doesn't exist across a network — the ack can always be lost. What I build is at-least-once delivery plus idempotent consumers keyed on a client-supplied ID, which is observationally exactly-once and is what every system claiming exactly-once actually does internally." |
| "Add a bigger queue." | "A queue doesn't create capacity, it converts a throughput problem into a latency problem. If the queue is persistently non-empty the pool is undersized. I'd bound the depth, shed above it, and autoscale on queue age, not depth." |
| "Isn't a central control plane a SPOF?" | "It's a deliberate one, and the mitigation is that it's not in the data path — if it's down, racks keep serving on their current config. I'd take that over a distributed control plane whose split-brain modes I can't reason about." |
| "Consistent hashing solves your sharding." | "It solves rebalancing — only 1/N of keys move on a membership change. It doesn't solve skew: one whale tenant still lands on one shard. Virtual nodes smooth the distribution; hot-key splitting is what handles the whale." |
| "You're over-engineering the rollout." | "The staged rollout is the only part I'd refuse to cut, because it's what bounds the blast radius of every future change. I'd cut the P2P distribution tier and the tail sampler first — both are optimizations with a clear trigger condition for adding them later." |
| Symptom | First metric | Usual cause | Fix |
|---|---|---|---|
| Queue depth rising 20 min | queue age, worker count | Undersized pool or a poison job | Scale out; dead-letter after N attempts. Adding queue depth makes it worse. |
| p99 up, CPU flat | queue wait, lock contention, GC | Queueing or locking, not compute | Bound concurrency, shed load |
| Cache hit rate collapse | key cardinality, TTL distribution | Key change on deploy, or synchronized TTL expiry | Jitter TTLs; single-flight on miss |
| Retry storm after a blip | upstream error rate + retry rate | Retries without jitter | Exponential backoff with jitter, circuit breaker |
| Storage growth 10× overnight | series-creation rate, ingest by source | High-cardinality label or a log-level change | Per-source quotas, drop-and-count |
| Fleet partially stale | observed-vs-desired version histogram | Push-based delivery missed offline nodes | Reconciliation loop, alarm on convergence lag |
| Rollback impossible | — | Rollback path depends on the thing that broke | Out-of-band recovery channel, last-known-good on the agent |
| Duplicate side effects | dedupe-table hit rate | At-least-once without idempotency | Idempotency keys on every mutating call |
W + R > N ⇒ strong consistency. Pick different values for different stores in the same design.Round 1: the coding screen · The AI system-design round · Gen-AI systems coding · The generic AI design framework
Problem breakdowns — Distributed rate limiter · Rate limiter, complete guide · Distributed job scheduler · System design framework and core concepts · Senior-level question bank
Why these six — AI200 Infrastructure Management Suite · AI200/AI250 rack-scale inference