End-to-End ML System Design Questions
Designing a complete machine learning system from problem to production. Covers the components and architecture of a production ML system, data flow from ingestion to serving, scalability, and integration of models into a larger product. Emphasizes the whole-system design tradeoffs that appear in ML system-design interviews.
A prototype that performed well in small-scale testing now needs to serve millions of users. Walk through how you would scale it up, and what you'd prioritize to avoid an embarrassing amount of downtime along the way.
Sample Answer
Direct answer
Priority order, not a shopping list: first decouple stateless request-serving from anything stateful (sessions, in-memory caches) so you can add replicas freely, second put an autoscaler in front of that stateless tier driven by a real load signal (queue depth or p95 latency, not just CPU), and third roll the whole thing out with staged traffic shifts (a canary) gated on service-level objectives (SLOs, the target thresholds for latency/availability you commit to) so a bad change is caught on 1% of traffic instead of 100%. Everything else (caching, async processing, monitoring) supports that spine.
Structured elaboration
1. Define the targets before touching infrastructure. Pin down p95/p99 latency (the 95th/99th percentile response time), an availability target (e.g. 99.9%), and a rough cost ceiling. Without these, "scale it up" has no stopping point and no way to know if a change helped.
2. Split the compute architecture into two paths.
- Synchronous, low-latency path: stateless model-serving replicas behind a load balancer, autoscaled on request-driven metrics (queue depth, in-flight requests, or p95 latency) rather than CPU alone, since CPU can look idle while requests queue on I/O.
- Asynchronous/batch path: anything that doesn't need an immediate response (bulk scoring, precomputation) goes through a durable queue consumed by an autoscaled worker pool, so a traffic spike on the async side doesn't compete with the latency-sensitive path for the same replicas.
3. Externalize state. Move sessions, feature lookups, and any per-request context out of the model-server process into a shared cache or store. This is what actually enables horizontal scaling: if state lives in the replica's memory, you can't add a second replica without splitting user traffic by session, which reintroduces the bottleneck you're trying to remove.
4. Add a caching layer for repeat/hot queries in front of the model-serving tier, sized to the fraction of traffic that's actually repeat, not universally, since caching stale predictions for a fast-moving model can itself be a correctness bug.
5. Observability before scale, not after. Golden signals (latency, traffic, errors, saturation) with alerting tied to the SLOs from step 1, so a slow rollout is visible before users report it.
6. Progressive rollout. Canary 1% of traffic, gated on automated SLO checks, then 5%, 25%, 100%, each stage paused until the previous stage's metrics are clean. Keep a one-click rollback path at every stage; this is what actually prevents "an embarrassing amount of downtime," not the capacity math itself.
Worked example
Say the prototype was validated at low volume, and the target is 2,000,000 daily active users, each issuing an average of 5 requests/day (a planning assumption, stated explicitly so the arithmetic is reproducible).
avg requests/day=2,000,000×5=10,000,000Convert to average queries per second (QPS, requests handled per second):
avg QPS=86,40010,000,000≈115.7Traffic isn't flat across the day; assume a peak-to-average factor of 3x (a common planning multiplier for consumer traffic, stated as an assumption here):
peak QPS≈115.7×3≈347Now assume load testing on a single replica measured a sustainable capacity of 20 QPS at the target p95 latency (this is the kind of number you'd get from your own load test, not a vendor benchmark, and it's the pinned input driving the rest of the math):
replicas needed=20347≈17.4→18 replicasAdd headroom for one replica's worth of failover (N+1) plus a burst buffer, say 30%:
18×1.3≈23.4→24 replicasSo the autoscaler's target ceiling for the synchronous serving tier is roughly 24 replicas at this projected peak, with the floor set by off-peak QPS using the same per-replica capacity figure. The point of doing this arithmetic explicitly is that it's re-runnable the moment your real load test gives you a different per-replica capacity number or your usage assumptions change.
Trade-offs & pitfalls
Over-provisioning for a peak factor you guessed wrong wastes real money every hour of every day; under-provisioning turns "millions of users" into an incident. Prefer measuring your actual peak-to-average ratio from prototype traffic over guessing, and re-derive the replica count once you have it. Stateful services (sticky sessions, in-memory model caches keyed by user) quietly block horizontal scaling even after you've "added autoscaling," so audit for hidden state before trusting the replica math. A canary only protects you if the signal it watches is fast and sensitive enough: SLO checks based on hourly aggregates won't catch a regression that matters within minutes. Finally, resist scaling complexity ahead of evidence: building a five-region, multi-tier architecture for a prototype that hasn't proven its growth curve yet is itself a way to introduce downtime, just earlier.
graph LR
A[Client request] --> B[Load balancer]
B --> C[Stateless model-serving replicas]
C --> D[Shared cache / session store]
B --> E[Async queue]
E --> F[Autoscaled worker pool]
C --> G[Monitoring: SLO dashboards]
G --> H[Canary gate]
H --> I[Traffic ramp: 1% to 100%]
Design a visual search system that lets a user search a large catalog with an image instead of text, under a tight latency budget and a limited hosting cost per request.
Sample Answer
Direct answer
Do the expensive work offline: precompute an embedding for every catalog item once, encode the query image only once per request, and let an approximate nearest-neighbor (ANN, a search structure that finds near-matches in sub-linear time by trading a small amount of recall for large speed and cost savings) index do the heavy lifting instead of comparing the query against the whole catalog. That combination is what makes a tight latency budget and a low per-request cost compatible with a large catalog.
Structured elaboration
Components and data flow:
- Image encoder (a convolutional or vision-transformer network, typically fine-tuned with a metric-learning objective such as triplet or contrastive loss) maps any image to a fixed-length embedding vector such that visually or semantically similar items land close together in that vector space.
- Offline catalog-indexing pipeline: run the same encoder over every catalog image (batch job, not per-request), store the resulting embeddings, and build an ANN index over them (e.g. a graph-based index like HNSW, hierarchical navigable small world, or a compressed index like IVF-PQ, inverted-file with product quantization).
- Online query path: user's image → preprocessing (resize/normalize) → one encoder forward pass → ANN lookup returning the top-K candidates → a re-ranking stage → results returned.
- Re-ranking (the ranking half of a candidate-gen-then-rank funnel): apply a heavier, more precise scoring model or metadata-based business rules to only the small candidate set K, not the whole catalog. This two-stage design is exactly what keeps per-request cost low: the expensive computation runs on K items, never on the full catalog.
- Incremental re-indexing: catalog contents change continuously, so the index-build pipeline supports incremental add/remove of vectors on a schedule (hourly/daily) rather than a full rebuild on every change.
- Caching: cache results for frequently-repeated or visually-similar queries to avoid redundant ANN lookups.
Key design decisions and why: the encoder only ever runs online on the single query image, never on the catalog; embedding dimensionality is chosen to balance recall against index size and lookup cost; the ANN algorithm and its compression settings (e.g. PQ code size) are chosen for a specific recall/latency/memory trade-off point rather than maximizing any one of the three; index updates are incremental so catalog freshness doesn't require a full, latency-impacting rebuild.
Worked example
Two reproducible calculations that show why the architecture is shaped this way (using explicit, stated assumptions, not vendor benchmarks).
Memory footprint of raw embeddings. Assume a catalog of 50,000,000 items with a 256-dimensional embedding stored as 4-byte floats:
50,000,000×256=12,800,000,000 floats 12,800,000,000×4 bytes=51,200,000,000 bytes≈51.2 GBEffect of compressing the index (product quantization). If each vector is compressed to a 32-byte code instead of stored as raw floats:
50,000,000×32 bytes=1,600,000,000 bytes≈1.6 GB 1.6 GB51.2 GB=32× reductionThat 32x reduction in memory is the direct, computable reason product-quantized indexes are attractive at catalog scale, at the cost of some retrieval recall you'd validate empirically against your own labeled query set.
Why an exhaustive scan doesn't work at this scale, qualitatively. Comparing a query against all 50,000,000 embeddings directly means 50,000,000 distance computations per query. A graph-based ANN index instead scales roughly with the logarithm of catalog size in the number of hops it needs to traverse:
log2(50,000,000)≈25.6This is an order-of-magnitude illustration of why ANN indexing, not exact search, is what makes catalog-scale visual search affordable per request; it is not a hard latency guarantee, since actual hop cost, recall target, and index tuning all still have to be measured against your own data and traffic.
Trade-offs & pitfalls
ANN retrieval is approximate by construction: recall, latency, and memory trade off against each other via index parameters (graph connectivity, PQ code size), and that trade-off point has to be validated against a labeled query set, not assumed. When the encoder itself is retrained (to improve quality), every catalog embedding is now stale and the entire index needs re-encoding and rebuilding, which is a real operational event to plan for, not a background non-event. A fixed re-ranking threshold can miss items that are semantically similar but visually distinct, or over-favor near-duplicates; this needs its own evaluation set. Catalog churn (constant adds/removes) has to be handled without a full reindex, which pushes toward index structures and pipelines that support incremental updates. Finally, unlike collaborative-filtering-style recommendation, a content-embedding visual search system is not blocked by cold-start on brand-new items, since the item's own image is sufficient signal, which is a genuine advantage of this approach worth calling out explicitly.
graph LR
subgraph Offline
C[Catalog images] --> E1[Image encoder]
E1 --> V[Embeddings]
V --> I[Build ANN index]
end
subgraph Online
Q[Query image] --> E2[Image encoder]
E2 --> L[ANN lookup: top-K]
I --> L
L --> R[Re-ranker]
R --> O[Results]
end
A set of inference endpoints needs autoscaling, where a fast CPU model and a much heavier GPU model have very different latency targets. What signals would you scale on, and where does the plan break down under a sudden traffic spike?
Sample Answer
Direct answer
Scale the CPU-model pool and the GPU-model pool on DIFFERENT signals, because their cost and latency profiles are fundamentally different: the CPU pool can scale reactively on request rate or CPU utilization since a new replica is ready in seconds, but the GPU pool must scale on a LEADING indicator, such as queue wait time or in-flight-requests-versus-target-concurrency, because a GPU replica's startup (image pull, driver init, model weights onto device memory) takes far longer, so waiting for GPU utilization to climb before scaling is already too late. Under a sudden spike, the plan breaks down specifically at the GPU pool: the spike arrives in seconds, scale-out plus cold start takes far longer, so the queue backs up and the tail-latency service-level objective (SLO) is breached before new capacity ever comes online.
Structured elaboration
| Aspect | CPU-model pool | GPU-model pool |
|---|---|---|
| Scaling signal | Requests per second (RPS) or CPU utilization, reactive | Queue wait time or in-flight requests versus a target concurrency, a leading indicator |
| Cold-start time | Seconds (container start) | Tens of seconds to minutes (image pull, driver init, weights onto device) |
| Batching lever | Rarely needed | Dynamic request batching raises per-replica throughput, which changes what "capacity" the scaling signal must account for |
| Typical failure mode under a spike | Brief queuing, self-heals quickly | Latency SLO breach persists until GPU replicas finish warming up |
Why raw GPU utilization is a poor scaling signal on its own. With request batching in place, a GPU can show high utilization even while individual requests are queuing, because a bigger batch keeps the device busy while each request's wait time to be included in a batch grows. Utilization is a lagging, batch-size-confounded signal; queue-wait-time or concurrency-based signals catch a growing backlog before the device itself even looks busy.
Batching and latency are the same trade-off, not two separate knobs. A dynamic batching window raises throughput per replica and lowers cost per request, but every millisecond spent waiting to fill a batch is added directly to that request's latency. The batch timeout has to be sized against the same p95/p99 (95th/99th percentile) latency target the autoscaler is defending; loosen one and you have effectively loosened the other.
Surviving a spike is mostly not an autoscaling-speed problem. Because GPU scale-out is inherently slow, the realistic mitigation is a warm minimum-replica floor sized to the largest plausible burst, and/or an overload fallback to a smaller, CPU-servable model, rather than trying to make the reactive loop faster. Predictable, scheduled scaling (for known traffic patterns like time-of-day) reduces reliance on reactive signals catching up in time, but does nothing for a genuinely unannounced spike.
flowchart LR
Client[Client traffic] --> Router[Router splits by model type]
Router --> CPUQ[CPU-model queue]
Router --> GPUQ[GPU-model queue]
CPUQ --> CPUPool[CPU pool: fast reactive autoscale on RPS]
GPUQ --> GPUBatcher[Dynamic batcher]
GPUBatcher --> GPUPool[GPU pool: scale on queue-wait time]
CPUPool --> Metrics[Metrics: p95 latency, queue depth]
GPUPool --> Metrics
Metrics --> Autoscaler[Autoscaler decision loop]
Autoscaler --> CPUPool
Autoscaler --> GPUPool
Worked example
Suppose steady state is 100 requests/sec (RPS), fully served by 10 GPU replicas at 10 RPS each, and a new GPU replica takes 90 seconds to become ready (cold start). A spike to 400 RPS (4x) arrives in under 5 seconds. The autoscaler now needs 40 replicas total, 30 more than exist, and those 30 take 90 seconds to warm up. During that 90-second window, requests keep arriving at the higher rate while capacity hasn't grown yet, so the excess accumulates as backlog:
excess arrivals during warm-up≈(400−100) req/s×90 s=27,000 requestsThose 27,000 requests are either queued (and blow the latency SLO for everyone behind them) or dropped, before capacity ever catches up. This is the concrete arithmetic behind why the plan "breaks down": the backlog scales linearly with cold-start time, so halving cold-start time only halves the backlog, it does not remove it. Only pre-provisioned headroom, or explicit backpressure/shedding once the queue passes a threshold, avoids the backlog entirely.
Trade-offs & pitfalls
A single blended scaling signal across both pools hides the GPU pool's slow cold start behind the CPU pool's fast one; the two pools have to be scaled independently. A large warm replica floor solves the spike problem but pays for idle GPU capacity around the clock, so its size is a genuine cost-versus-risk decision, not a free win. The common wrong turn is chasing a "better" reactive metric (switching from CPU utilization to GPU utilization, for instance) when the real fix is a leading indicator plus pre-provisioned headroom; no purely reactive signal survives a cold-start-dominated spike. Backpressure, deliberately rejecting or degrading some requests once the queue passes a threshold, protects the SLO for the requests that are admitted, at the explicit cost of dropping others; a capacity-only plan often forgets to make that trade on purpose.
After a blue/green deployment, you discover that traffic on the new (blue) side is producing subtly biased results because of a small mismatch in how data was preprocessed between staging and production. What would you put in your testing and validation process to have caught this before it shipped?
Sample Answer
Direct answer
The gap that let this ship is a validation process that checked the model's outputs but never directly compared the staging and production feature pipelines against each other on the same inputs. The fix is to add an explicit parity check, a statistical test that compares the distribution of every feature as it lands in production against the distribution seen in staging (or training), gated as a hard blocker before blue traffic is ramped, not an optional dashboard someone glances at after the fact.
Structured elaboration
Where the parity check sits in the pipeline
flowchart LR
A[Training data] --> B[Preprocessing spec v1: versioned and hashed]
B --> C[Staging pipeline]
B --> D[Production pipeline]
C --> E[Feature distribution sample: staging]
D --> F[Feature distribution sample: prod]
E --> G[PSI distribution-diff test]
F --> G
G --> H{PSI within threshold}
H -->|No| I[Block blue rollout]
H -->|Yes| J[Shadow traffic on blue]
J --> K[Canary ramp with rollback gate]
1. Pipeline parity, verified, not assumed
- Preprocessing logic (scalers, encoders, tokenizers, normalization constants) has to be a single versioned artifact loaded identically by staging and production, not two independently maintained code paths that happen to be intended to match.
- Even with a shared artifact, a parity test still matters: run the same batch of real (or replayed) inputs through both environments and diff the outputs field-by-field. A silent mismatch (log1p applied in one place and log10 in another, a timezone offset in a time-based feature, a different null-fill value) shows up as a diff here even when both pipelines "look correct" individually.
2. Distribution-diff testing as an automated gate
This is the check that catches the class of bug in this scenario: nothing crashed, no schema changed, but the feature values are subtly on a different scale. Bucket each feature into bins and compare the proportion of production traffic landing in each bin against the expected (staging or training) distribution using the population stability index (PSI), a standard measure of how much a distribution has shifted:
where ai is the actual (production) proportion in bin i and ei is the expected (staging) proportion. A PSI above roughly 0.2 is the common industry rule of thumb for "this is a material shift, not noise" and should block promotion.
3. Where this sits in the deployment pipeline
- Schema and contract tests (types, ranges, required fields) run first in CI, on every change, and catch structural breaks.
- The distribution-diff test runs against a production-like traffic sample before blue gets any real traffic, and again continuously once blue is in shadow mode, comparing shadow predictions and their input features against the green baseline on the same live traffic.
- Shadow mode: route a copy of real production traffic through blue without acting on its output, and compare blue's predictions and confidence distribution against green's on the same requests. A processing mismatch that changes the input distribution will usually show up as a shift in blue's output distribution too, not just its accuracy on a later-arriving label.
- Canary ramp (a few percent of real traffic) with an automatic rollback gate tied to the same distribution-diff and bias metrics, not just latency and error rate.
4. Governance around the pipeline itself
- A pre-deploy checklist with explicit sign-off from whoever owns the data/feature pipeline, separate from whoever owns the model, since this bug sits exactly at the seam between those two areas of ownership.
- An automated diff tool that flags any change to normalization constants, encoders, or tokenizer vocabulary as a reviewed, called-out change, not a side effect buried in an unrelated pull request.
Worked example
Suppose a feature (say, a scaled transaction amount) has this expected (staging/training) distribution across four bins, and this is what's actually observed in production after the scaling mismatch:
| Bin | Expected (staging) | Actual (production) |
|---|---|---|
| Low | 0.10 | 0.05 |
| Medium | 0.40 | 0.25 |
| High | 0.35 | 0.40 |
| Very high | 0.15 | 0.30 |
A PSI of 0.216 clears the ~0.2 "material shift" threshold, which is exactly the kind of quiet mass-shift toward the "very high" bin a scaling mismatch (for example, a log1p transform in staging versus a log10 transform in production) produces. Wired into the promotion pipeline as a hard gate, this catches the bug before blue takes real traffic, instead of after clinicians, users, or downstream consumers see biased output.
Trade-offs & pitfalls
- Schema tests alone are not enough: this bug passed every type and range check because nothing was structurally wrong, only the values were subtly rescaled. The distribution-diff test is the piece that closes that gap, and it's easy to skip because it takes real engineering effort to define good bins and thresholds per feature.
- Setting the PSI (or equivalent) threshold too loose defeats the purpose; setting it too tight creates alert fatigue and teams start ignoring it, which is its own failure mode. The threshold needs to be tuned per feature against historical natural variation, not copy-pasted as a single global number.
- Comparing distributions once at deploy time and never again misses drift that develops after a clean launch; the same test needs to run continuously as a monitoring signal, not just as a pre-deploy gate.
- Bias specifically (as opposed to a generic accuracy regression) requires checking the diff broken out by subgroup, not just in aggregate, since a shift that is invisible in the pooled distribution can be concentrated in one subgroup.
- Rollback has to be automatic and fast (traffic-weight based, not a redeploy), or the gate finding the problem doesn't actually limit the blast radius.
You're asked to sign off on a new model feature before it goes fully live. What would actually be on your go/no-go checklist, spanning security, reliability, cost, and compliance, and what would make you say no even if the model's accuracy looks fine?
Sample Answer
Direct answer
A passing accuracy number is necessary but never sufficient. Sign-off is really a check on whether the whole SYSTEM, not just the model, is safe to expose: security, reliability, cost, and compliance each get their own concrete checks, and any single hard-fail among them (no rollback path, unreviewed handling of personal data, unbounded cost exposure, or an unmet regulatory requirement) is grounds to say no regardless of how good the accuracy looks.
Structured elaboration
| Category | Concrete checks |
|---|---|
| Security | Least-privilege access to the model endpoint and training artifacts; no secrets embedded in code or images; input validation against malformed or adversarial payloads; dependency and container-image scanning on the serving image |
| Reliability | Defined service-level objectives (SLOs, e.g. a target p95 latency and error rate) tested under expected AND above-expected load; a rollback path that has actually been exercised, not just documented; monitoring and alerting wired up and firing correctly BEFORE go-live; a named on-call owner |
| Cost | A unit-economics estimate (cost per 1,000 inferences) with an explicit budget ceiling; an autoscaling ceiling so a traffic spike or bug produces a bounded bill, not an unbounded one |
| Compliance | A valid, checked basis for any personal data used in training or inference; for a regulated decision, a verified (not assumed) explainability and audit story; a retention and deletion policy for logged inputs and outputs that matches actual data-handling commitments |
Hard-no triggers even when accuracy looks fine:
- No working rollback: a bad deploy can't be undone quickly, so the actual risk being signed off on (what happens when this breaks) is unaddressed.
- The validation population doesn't match production traffic, a training-serving skew risk that a good offline number can mask entirely.
- No monitoring wired up: a live regression would be invisible until someone notices downstream damage.
- A compliance gap on a regulated use case (for example, a decision affecting people's credit or employment) without the required documentation or basis.
- Cost with no ceiling: a single bug or traffic spike could produce an unbounded bill.
Each checklist item needs a piece of verification evidence attached, a test result, a dashboard link, a signed-off document, not just a checkbox; a checklist without evidence degrades into a rubber stamp.
Worked example
A recommendation model shows a 4% offline lift and clears every accuracy gate. During review it's found the serving endpoint has no automated rollback (reverting to the previous version would require a manual, multi-step redeploy) and there's no alert on the serving error rate. The decision is NO-GO until automated rollback exists and the error-rate alert is wired and tested, regardless of the 4% accuracy number, because what's actually being signed off on is "what happens in the first minutes after this breaks in production," and today the honest answer is "nobody would know, and fixing it would be slow."
Trade-offs & pitfalls
The single most common mistake in this review is over-indexing on accuracy, since it's the number that's already been validated and the hardest one to say no to once it looks good. A checklist with no owner assigned per item tends to degrade into theater over time; assign accountability per item, not just per category. Treating compliance as "someone else's problem to check later" is exactly the item most likely to block a launch late and expensively, after infrastructure and go-live plans are already committed, if it isn't checked at sign-off instead.
Unlock Full Question Bank
Get access to all End-to-End ML System Design interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.