← Interview Mastery
IC4IC5IC6

Qualcomm Round 2: Classic System Design

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.

45 min read · 15 sections
0

1. Quick anchor

There are two different design rounds in this loop and candidates routinely prepare for only one of them.

  • The AI design round — "design a RAG system," "design a multi-tenant inference service." That's the Gen-AI system-design primer.
  • This round — ordinary distributed systems, graded at staff. Rate limiters, queues, schedulers, metrics pipelines, control planes. No model in sight.

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

2. The drive script

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.

  • "Let me size this before I pick a database."
  • "I'll take availability over consistency here, and here's the specific anomaly the user will see because of it."
  • "That single point of failure is deliberate — here's why, and here's the fallback when it's down."

3. The numbers you carry in

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

  • 1 Gbps = 125 MB/s; 10 Gbps = 1.25 GB/s.
  • One day ≈ 86,400 s ≈ 10⁵ s. So 1M events/day ≈ 12/s; 100M/day ≈ 1,200/s; 1B/day ≈ 12,000/s.
  • Peak ≈ 2–3× average. Design for peak, bill for average.
  • A single modern server: ~10–50K simple QPS, ~1–5K QPS with real work.
  • Redis: ~100K ops/s per instance, sub-millisecond.
  • One TB/day = ~12 MB/s sustained. Say it that way — it makes "we'll just write it all" sound as expensive as it is.

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.

4. The primitive toolbox

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.

5. Six worked designs

D1 — Distributed rate limiter

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

  1. Centralized store (Redis) with an atomic script. The whole read-modify-write must be one atomic operation or 40 nodes race and every one of them sees budget. Correct, ~0.3 ms per request, and a hard dependency.
  2. Local buckets at 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.
  3. Local buckets with periodic reconciliation. Each node keeps a local allowance and syncs usage every second, redistributing unused budget. Near-zero hot-path latency with bounded error. This is the staff answer, and the thing that makes it staff is stating the error bound: "a tenant can exceed the limit by at most (nodes × sync interval × rate) in the worst case — about 3% at our numbers."
# 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.

  • Sharding + hot keys. Shard buckets across Redis instances by consistent hash of the key. Then the failure it doesn't solve: one enormous tenant lands on one shard. Fix by splitting that key into K sub-buckets at rate/K, chosen by hashing the request ID.
  • What happens when Redis is down. This is the question. Fail open — a rate limiter that takes down the API it protects has inverted its own purpose — but fail open onto a conservative local limiter, not onto nothing. Then alarm loudly, because you are now unprotected against the abuse case.

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.


D2 — Rack-scale telemetry pipeline

"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,000

Then 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 + dashboards

The five decisions worth defending.

  1. Push vs pull. Pull (Prometheus-style) gives the collector control of the rate and makes "is it up?" free — but needs service discovery for 80,000 targets and struggles across NAT/customer boundaries. Push scales to on-prem appliances and short-lived jobs but lets a misbehaving agent flood you. Choose per environment: pull inside your own DC, push from customer appliances, with a per-source quota on the push path.
  2. Kafka in the middle, always. Without a buffer, a TSDB hiccup becomes permanent data loss and back-pressures 80,000 agents. With it, the TSDB can be down for an hour and you replay.
  3. Retention tiers, priced. Raw 1 s for 15 days → 1 min rollups for 90 days → 1 hour for 2 years. A 60× reduction at the first hop is where the storage bill actually gets decided.
  4. The alerting path must not share the query path. If alerts are evaluated by the same service that serves dashboards, a dashboard query storm suppresses your alerts precisely when things are on fire. Separate services, separate capacity.
  5. Aggregate at the edge. A rack-local agent that pre-aggregates 64 devices' metrics cuts the fan-in 64×. The tradeoff: you lose per-device granularity for anything you roll up, so keep per-device for the handful of signals that drive replacement decisions (thermals, ECC errors, throttling).

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.


D3 — Distributed job scheduler

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

  • The lease is what makes worker death survivable. No heartbeat → the lease expires → another worker reclaims it. No separate failure detector needed.
  • attempts bounds retries. Past MAX_ATTEMPTS, move to a dead-letter queue with the error — never retry forever, and never silently drop.
  • The 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.
  • At-least-once + idempotent execution. The lease can expire while the worker is still alive but partitioned, so a job will occasionally run twice. That's not a bug to eliminate, it's a contract to design for — job handlers take an idempotency key derived from (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.


D4 — Fleet control plane: config and rollout

"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 + health

Four properties, four sentences.

  1. Config is an immutable, versioned, signed artifact. Never mutable rows. This gives you diff, rollback, audit, and integrity in one decision.
  2. Staged rollout with health gates. 1 rack → 1% → 10% → 50% → fleet, with a bake time and an automatic rollback trigger at each gate. The wave caps live in the control plane, not in the operator's head — that's what makes them a guarantee instead of a convention.
  3. Agents keep last-known-good. If the new config fails a local health check, revert and report unhealthy. This is what turns a bad push into a metric instead of an outage.
  4. The control plane must not be in the data path. If it's down, racks keep running on their current config — degraded management, not degraded service. Say this explicitly; it's the single most important property and interviewers wait for it.

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.


D5 — Model artifact distribution

"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 days

Say that out loud and the design writes itself.

The fix, in three layers.

  1. Content-addressed, chunked storage. Split artifacts into content-defined chunks and address every chunk by its hash. You get deduplication, integrity verification, and resumable transfer from one decision. Consequence worth stating: two model versions that share 80% of their weights transfer only the 20% that differs, so an incremental release is ~10 GB, not 50.
  2. Hierarchical distribution. Central store → one cache per rack (or per rack row) → devices over the LAN. Now WAN transfer is 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.
  3. Warm-ahead staging. Pre-position the next version to rack caches before the rollout window, so the cutover is a local copy and a restart, not a download. This is what makes rollout time predictable, and predictability is what lets D4's health gates have meaningful bake times.

Two more things to name.

  • Signing and verification are non-negotiable. An artifact store that distributes unsigned binaries to 80,000 accelerators is the highest-value supply-chain target in the company. Verify the signature on the agent, before load, every time — not just at upload.
  • Garbage collection. Old versions accumulate at every cache tier and will silently fill every disk in the fleet. Refcount by "versions any rack might roll back to," keep N-2, and alarm on cache disk usage — this is the boring failure that actually happens.

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.


D6 — Log and trace storage

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

6. Pushback drills

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

7. Failure-mode catalogue

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

8. Flashcards

  • Requirements → estimates → API → data model → high level → deep dive → bottlenecks.
  • 1 day ≈ 10⁵ s. 1M/day ≈ 12 QPS. 1B/day ≈ 12K QPS. Peak ≈ 2–3× average.
  • RAM 100 ns · NVMe 50 µs · same-DC RTT 0.5 ms · cross-region 100 ms.
  • 1 Gbps = 125 MB/s. 1 TB/day = 12 MB/s sustained.
  • W + R > N ⇒ strong consistency. Pick different values for different stores in the same design.
  • Consistent hashing fixes rebalancing, not skew. Virtual nodes for smoothing; key-splitting for whales.
  • Exactly-once = at-least-once + idempotent consumer. Always.
  • A persistently non-empty queue is a mis-sized system.
  • Retries always need jitter.
  • Control planes reconcile; they don't push. And they stay out of the data path.
  • Cardinality kills a metrics system before volume does. Alarm on the rate of series creation.
  • Tail sampling keeps 100% of errors and 1% of the boring traces.
  • Every rollout needs a rollback path that doesn't depend on the rollout path.

9. Further reading

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

Primary sources
← More in Interview Mastery