ML System Design · Recommendation · Interactive

Designing a Recommender for 100 Million Users

The full ML system design interview, worked end to end: a video home feed at 28,000 requests per second. Every stage, every number, and the follow-up questions that separate a senior answer from a staff one.

The short version

A scalable recommender is a funnel. Each stage trades accuracy for throughput, and the whole design is deciding where to spend compute.

“Design the home feed for a video streaming service.” It is the most common ML system design prompt there is, and it is common for a reason: it touches every hard part of production ML at once. Data you have to log correctly. Labels that lie. A latency budget measured in milliseconds. A catalog that changes daily. Feedback loops that train the model on its own output.

This post works the problem the way a strong interview answer does: in order, out loud, with numbers. Each section ends with the follow-up an interviewer is likely to push on.

01Clarify before you design

Spend the first five minutes on questions. Each one should change a design decision. If an answer would not change anything, do not ask it.

QuestionAssumeWhat it changes
What surface?Home page: rows of titlesPage-level layout, not a single ranked list
Business model?Ad-supported streamingWatch time pays the bills; clicks do not
Scale?100M DAU, 300M users totalForces a multi-stage funnel
Catalog?~5M items, ~10K new per dayCold start is a first-class problem
Latency?p99 under 200 msHeavy model only on a few hundred items
Logged in?Mix of signed-in and anonymousNeed a session-only fallback path
Constraints?Kids profiles, licensing by regionHard filters before ranking, never after

Then say the scope out loud. “I’ll design personalized home-feed ranking. I’ll treat search, ads insertion, and thumbnails as out of scope unless we have time.”

› Interviewer follow-up: “Why not just ask for all the requirements up front?”

Because the goal is to show judgment, not to collect a spec. A strong candidate asks the three or four questions that fork the design, states assumptions for the rest, and moves on.

A useful test: after each answer, say which component it affected. “Kids profiles means a hard filter stage before retrieval results reach the ranker.” That sentence is worth more than five extra questions.

02Turn the business goal into a label

This is where most answers are weakest. There are four layers, and each is a proxy for the one above it.

LayerExampleWhy it is not enough
Business goalLong-term retention, revenueMoves slowly. You cannot train on it.
Online metricWatch time per DAU, 7-day return rateMeasured only in A/B tests
Offline metricNDCG, recall@K, AUC, calibrationEvaluated on logged data from the old policy
Training labelPer impression: click, watch seconds, completionEach label rewards a different behavior

Pick the label deliberately

The answer is all of them, predicted as separate heads, then combined. That is section 08.

The line interviewers remember

“I won’t optimize clicks alone, because that trains the model to be a thumbnail classifier. I’ll predict several engagement signals and combine them with weights we tune against long-term retention in A/B tests.”

› Go deeper: how to predict watch time with a classifier

YouTube’s 2016 paper used weighted logistic regression. Positive (clicked) impressions are weighted by watch time; negatives get weight 1. The learned odds then approximate expected watch time per impression, so at serving you rank by elogit.

Why not plain regression on seconds? The distribution is heavy-tailed and zero-inflated. Most impressions are not clicked. A classification framing handles the zeros naturally, and the weighting handles the magnitude.

A modern alternative: predict p(click) and E[watch | click] as two heads, and multiply. This makes each head easier to calibrate and debug.

› Interviewer follow-up: “Your label is biased by what you showed. How do you handle that?”

Two biases matter. Position bias: items in the top-left get clicked because of where they are. Exposure bias: items never shown have no labels at all.

For position, train with position as an input feature (or a separate shallow “bias tower”) and set it to a fixed value at serving. YouTube’s 2019 ranking paper does exactly this with a shallow tower added to the main logit.

For exposure, reserve a small slice of traffic for exploration (randomized or bandit-chosen slots). It gives you unbiased data for evaluation and for new items. Log the propensity of each impression so you can reweight later.

03Do the arithmetic

Numbers turn a hand-wavy architecture into a forced one. The calculator below runs the same math I would do on a whiteboard. Change any assumption and every downstream number updates.

Four numbers from that calculator drive the rest of the design.

› Interviewer follow-up: “Your item feature lookups are 28 million per second. What breaks?”

A remote key-value store at that rate is expensive and adds tail latency. But the item set is only 5M rows, and item features change slowly.

So cache item features in-process on the ranking servers, refreshed every few minutes from the feature store. Only user and real-time session features go over the network, which is one or two lookups per request, not a thousand.

This is the general pattern: find the side of the join that is small and slow-changing, and replicate it.

04The architecture: a funnel in three time scales

The system has three paths that run at different speeds. Draw all three. Many candidates only draw the online path, and the interviewer then asks where the embeddings come from.

Online path (milliseconds) serves requests. Nearline (seconds to minutes) keeps features fresh. Offline (hours) trains models and builds the index.

The funnel

StageItems in → outModelCost per item
Retrieval5M → ~1,000Several generators; two-tower + ANN is the main one~nothing (index lookup)
Filtering1,000 → ~900Rules: region, maturity, already watchedMicroseconds
Light ranking900 → 300Small MLP or GBDT, distilled from the heavy ranker~0.1 MFLOP
Heavy ranking300 → 300 scoredMulti-task deep model with sequence features~50 MFLOP
Re-ranking300 → ~50 on the pageDiversity, freshness, business rules, explorationList-level

Each stage exists because the next one is too expensive to run on its input. That is the whole argument, and saying it explicitly is better than listing stages.

The latency budget

A p99 of 200 ms sounds generous until you split it. Candidate generators run in parallel, so their cost is the slowest one, not the sum.

Illustrative p99 budget. The heavy ranker gets the largest share. The slack line absorbs tail variance and retries.
› Interviewer follow-up: “Why not precompute everyone’s feed offline?”

You can, and for some products it is the right call. Precomputing daily feeds for 300M users is a large but simple batch job, and serving becomes a key-value read.

What you lose is session context. If I just watched three cooking shows, the feed should reflect that on my next page load, not tomorrow. You also waste compute on users who never come back that day.

A good hybrid: precompute a candidate pool per user offline (cheap, broad), and do ranking online with fresh session features. Fall back to the precomputed list if online ranking times out.

05Data and logging: where most systems break

The model is only as good as the join between what you showed and what happened.

Log the impression, with its features

Every request gets a request_id. For every item shown, log the position, the model version, the propensity, and the exact feature values used to score it.

Then join engagement events (clicks, watch progress, likes) back to impressions by request_id and item.

Why log features instead of recomputing them

If you recompute features at training time, you get today’s value for an event that happened last week. The model learns from information it will never have at serving. This is training/serving skew, and it produces models that look great offline and flat online.

Logging serving-time features makes training data correct by construction. It costs storage. It is worth it.

Handle delayed and missing labels

› Go deeper: correcting calibration after negative downsampling

If you keep a fraction w of negatives, the model’s predicted odds are inflated by 1/w. Correct at serving with p = q / (q + (1 − q) / w), where q is the raw prediction.

Calibration matters here more than in most classifiers, because the value function multiplies probabilities from several heads. A head that is off by 2× silently doubles its weight in the final score.

06Features and the feature store

GroupExamplesFreshness
User, long-termGenre affinities, embedding, tenure, device mixDaily batch
User, sessionLast 50 watched items, time since last visit, current rowSeconds (streaming)
ItemContent embedding, genre, length, age, languageDaily, plus on upload
Item, popularityViews last hour / day, completion rate, trend slopeMinutes
ContextTime of day, device, country, networkPer request
User × itemWatched this series before, affinity to this genre, similarity to recent watchesComputed at request

The user × item cross features usually carry the most signal. They are also the most expensive, because they are computed per candidate.

A feature store gives you one definition of each feature, used for both offline training and online serving. Its key job is point-in-time correctness: when building training data, each row sees feature values as of the event time, never later.

› Interviewer follow-up: “How do you represent a user’s watch history?”

Three options, in increasing cost. Average the embeddings of recent items (cheap, loses order). Attend over history using the candidate as the query, as in DIN (target attention). Or run a transformer over the sequence, as in BST and most current systems.

Target attention is usually the best value. The relevant part of my history depends on what you are scoring. My cooking watches matter for a cooking show, not for a thriller.

Cap the sequence length (50–200 items) for latency, and precompute the candidate-independent part once per request, not once per item.

07Candidate generation

Use several generators, not one. Each covers a failure of the others.

GeneratorWhat it catchesShare
Two-tower + ANNPersonalized relevance, the main source~50%
Item-to-item (co-watch)“Because you watched X”~20%
Continue watchingUnfinished series and moviesAll eligible
Trending / popular by regionNew users, current events~10%
Fresh contentNew uploads with little data~10% (exploration budget)
Followed / subscribedExplicit interestAll eligible

The two-tower model

A user tower maps user and context features to a vector. An item tower maps item features to a vector. The score is their dot product.

That constraint is the point. Because users and items never interact until the dot product, item vectors can be computed offline and indexed. At request time you compute one user vector and do a nearest-neighbor search.

Training the two-tower model# in-batch softmax with sampling-bias (logQ) correction u = user_tower(user_feats) # [B, d], L2-normalized v = item_tower(item_feats) # [B, d], L2-normalized logits = (u @ v.T) / temperature # [B, B]; diagonal = positives logits -= log(item_frequency)[None] # popular items appear as negatives more often loss = cross_entropy(logits, arange(B))

The index

At 5M items, HNSW on one machine is fast and simple, with recall above 95% at sub-millisecond latency. At hundreds of millions of items, use IVF with product quantization, or ScaNN, and shard.

Rebuild the index when the item tower retrains. Between rebuilds, insert new items incrementally. Keep the user and item towers version-locked: a user vector from model v7 against an index from v6 returns garbage, and nothing errors.

For the full treatment of this co-design, see Two Towers, One Index.

› Interviewer follow-up: “How do you evaluate retrieval separately from ranking?”

Use recall@K: of the items the user actually engaged with next, what fraction appear in the top K retrieved? Retrieval’s job is to not miss things. Precision is the ranker’s job.

Measure it per generator and for the merged set. A generator whose items are never in the final page, and never add recall, is costing latency for nothing.

Watch out: recall against logged engagement only counts items the old system showed. Evaluate on exploration traffic too.

08Ranking: multi-task, calibrated, combined

The heavy ranker sees a few hundred candidates, so it can afford rich features and cross interactions.

Architecture

The value function

The final score combines the heads. Keep it outside the model.

Value model, applied after the rankerscore = w1 * p_click * E[watch_sec] + w2 * p_complete + w3 * p_like - w4 * p_not_interested

Product teams can shift priorities by changing weights, without retraining. You tune the weights with A/B tests against the long-term metric. This only works if every head is calibrated.

› Interviewer follow-up: “Why not one model per task?”

Serving cost and data efficiency. Five models means five forward passes on 300 items at 28K QPS. Sharing a bottom lets sparse tasks (likes) borrow representations learned from dense tasks (clicks).

The risk is negative transfer: one task’s gradient hurts another. MMoE-style gating and per-task loss weights manage it. If a task still degrades, split it into its own tower or model.

› Go deeper: the light ranker and distillation

The light ranker must be cheap enough for ~1,000 items and agree with the heavy ranker on what goes in the top 300. Train it to match the heavy ranker’s scores (distillation), not just the labels.

Measure it as a filter: of the heavy ranker’s top 50, what fraction survives the light ranker’s cut? If that drops, the heavy ranker never sees its best items, and no heavy-ranker improvement will show up online.

09Re-ranking: the page, not the list

Sorting by score gives you ten episodes of the same show. Re-ranking fixes the page.

› Interviewer follow-up: “Diversity lowers your offline NDCG. How do you justify it?”

Offline NDCG scores each item independently and assumes the user wants the top item most. It cannot see that the fifth episode of a show is worth less once the first four are on the page.

Diversity is justified online: sessions get longer, and users who see a broader mix retain better. Say this, and propose measuring it with an A/B test plus a coverage metric (fraction of catalog shown per week).

10Cold start

New users

New items

11Training at scale

Embedding tables dominate

Dense layers are megabytes. Embedding tables for 300M users and 5M items are hundreds of gigabytes. The two need different parallelism.

Retraining cadence

Freshness has diminishing returns. Measure it: train on data that is one hour, one day, and one week old, and compare. The curve tells you how much pipeline to build.

› Interviewer follow-up: “Your model trains on data produced by your model. What goes wrong?”

A feedback loop. Items the model likes get shown, collect engagement, and become more liked. Items it ignores never get data. Over time the catalog coverage shrinks and the model gets confident about a smaller world.

Mitigations: the exploration slice, logQ and popularity corrections, propensity-weighted training on the exploration data, and a coverage metric on the dashboard. Also keep a small holdout of users served by a simple non-personalized policy, as a reference point that the loop cannot contaminate.

12Serving at scale

13Evaluation: offline chooses, online decides

Offline

Online

Why offline and online disagree

Offline metrics are computed on data chosen by the old model, one item at a time, without the page around it. Online, the new model changes what users see, the page is re-ranked, and users adapt. Treat offline gains as permission to run an A/B test, not as a result.

14Monitoring and failure modes

FailureSymptomDetect with
Training/serving skewOffline win, online flatCompare logged vs. recomputed feature distributions
Broken upstream featureSudden metric dropNull rate and range checks per feature, per hour
Tower version mismatchRelevance collapses, no errorsVersion tag on every vector; assert match at query time
Stale indexNew items missing from feedAge of newest item in index
Feedback loopCoverage shrinks over weeksWeekly catalog coverage and Gini of impressions
Calibration driftValue function misweights headsPredicted vs. observed rate per head, daily
Latency regressionp99 creeps up, timeouts risep99 per stage; fallback-serve rate

The fallback-serve rate deserves its own alert. When the heavy ranker times out, users still get a page, so nothing looks broken. Quality quietly drops.

15The follow-ups that separate levels

Senior answers name the trade-off. Staff answers say how they would measure it and what would change the decision.

QuestionStrong answer, in one line
How many candidates should retrieval return?Sweep K; stop where recall@K of the final page flattens and ranker latency still fits.
GBDT or deep model for ranking?GBDT to start and for dense features. Deep once sparse IDs and sequences dominate.
How fresh must the model be?Measure the accuracy-vs-staleness curve, then buy only the freshness it pays for.
What if watch time goes up but retention goes down?The value weights are wrong. Re-tune against the long-term holdout, add a satisfaction head.
How do you debug a bad recommendation?Trace it: which generator, which light-ranker score, which heads, which re-rank rule.
How do you scale to 10×?Shard the index, cache more, shrink the heavy ranker’s input, not its accuracy.
Where would you use an LLM?Offline: item understanding, content embeddings, cold-start descriptions. Not in the 50 ms ranking path.

16Run the clock

The most common failure is not a wrong answer. It is running out of time before the evaluation and monitoring sections, which are where depth shows. Budget a 45-minute round like this:

Leave evaluation and serving a protected slot. Interviewers score what they heard, and they cannot score the back half if you never reached it.

Takeaways

  1. Objective first. The label decides the behavior. Predict several signals and combine them outside the model.
  2. Numbers force the architecture. 28K QPS and 8M scored items per second is why there are four stages.
  3. Log serving-time features. It is the cheapest insurance against skew.
  4. Calibrate every head. The value function is only as good as the probabilities it multiplies.
  5. Protect exploration. Without it, cold start and feedback loops eat the catalog.
  6. Offline chooses, online decides. Guardrails and a long-term holdout keep you honest.

References

Covington et al., Deep Neural Networks for YouTube Recommendations, RecSys 2016 · Zhao et al., Recommending What Video to Watch Next, RecSys 2019 · Yi et al., Sampling-Bias-Corrected Neural Modeling for Large Corpus Item Recommendations, RecSys 2019 · Ma et al., MMoE, KDD 2018 · Tang et al., PLE, RecSys 2020 · Wang et al., DCN V2, WWW 2021 · Zhou et al., Deep Interest Network, KDD 2018 · Naumov et al., DLRM, 2019 · Liu et al., Monolith, 2022 · Wilhelm et al., Practical Diversified Recommendations on YouTube with DPPs, CIKM 2018 · Malkov & Yashunin, HNSW, 2016 · Deng et al., CUPED, WSDM 2013 · Netflix Tech Blog, Innovating Faster on Personalization Algorithms at Netflix Using Interleaving, 2017

The capacity calculator is plain arithmetic on the assumptions shown in it. The latency budget, per-item FLOP costs, and generator shares are illustrative round numbers for a system of this size, not measurements of any company’s production system.

Read next: Two Towers, One Index · Why Did You Rank That?

Cite this post
@misc{murugesan2026recsys,
  author = {Murugesan, Sugeerth},
  title  = {Designing a Recommender for 100 Million Users},
  year   = {2026},
  month  = {sep},
  url    = {https://sugeerth.github.io/blog/recsys-system-design/},
  note   = {Accessed: [date]}
}
SM
Sugeerth Murugesan Staff ML Engineer / Scientist · Intel / Intuit