A recommender answers one question many times a second: out of everything we could show this person right now, which twenty things should we show, in what order? The catalogue may hold millions of items, the latency budget is a few hundred milliseconds at most, and the only signal about what worked is what users did with what you showed them last time. No single model can score millions of items per request within that budget, so every serious system is a funnel of stages, each cheaper per item than the next and each seeing fewer items.
This article builds that funnel from first principles. It covers where candidates come from, how a two-tower retrieval model is trained and served, what the ranker predicts and how its outputs are combined, why re-ranking exists, and the data loop that feeds training. A worked example walks one request through the system, and the article ends with failure modes, trade-offs and a checklist. Approximate nearest-neighbour search itself is covered in HNSW, and feature serving in feature store architecture.
The constraint that shapes everything
Start with arithmetic. Suppose a video catalogue holds 10 million items and the ranking model costs about 20 microseconds of CPU per item scored. Scoring every item for one request would take 200 seconds of CPU, which no 150-millisecond budget can absorb.
The answer is to spend compute unevenly. A cheap stage narrows millions to a few thousand. A more expensive model scores those thousands. A final stage reorders the top few hundred with rules and diversity logic that need the whole list in view. Each stage can afford a richer model because it sees fewer items. Recall is decided early: an item that candidate generation misses can never be shown, however good the ranker is. Precision is decided late.
Candidate generation: many cheap sources
Candidate generation does not have to be one model. Most systems merge several sources, each good at a different kind of recall, and give each a quota so no single source dominates.
| Source | What it finds | Typical mechanism |
|---|---|---|
| Embedding retrieval | Items similar to the user's taste in general | Two-tower model, user vector queried against an ANN index of item vectors |
| Co-visitation | Items people consumed right after the ones this user just consumed | Precomputed item-to-item counts over sessions, looked up for recent items |
| Social or follow graph | New items from creators or accounts the user follows | Inverted index from creator to recent items |
| Trending and fresh | Items that are new or rising, which other sources have little data for | Time-decayed popularity per region or language |
| Exploration | Items the system is unsure about | Random or bandit-chosen sample from an eligible pool |
After the sources return, a merge step removes duplicates, drops items the user has already seen or hidden, applies hard policy filters such as region availability, and truncates to the ranker's budget. Filtering before ranking saves compute; filter again at the end, because re-ranking can pull items from a fallback list.
The two-tower retrieval model
The workhorse source is a two-tower model. One network, the user tower, turns user features such as recent watch history, language and device into a vector. A second network, the item tower, turns item features such as text, category, creator and age into a vector of the same size. The score is the dot product of the two. Because the towers never see each other's inputs, item vectors can be computed offline for the whole catalogue and loaded into an approximate nearest-neighbour index, while the user vector is computed once per request. Retrieval becomes a top-k inner-product search, which an index such as HNSW answers in milliseconds.
The usual training objective treats each logged positive pair, a user and an item they engaged with, as a classification problem over the other items in the same batch. Every other item in the batch acts as a negative. Popular items appear in many batches, so they are over-represented as negatives and get pushed down unfairly; the standard correction subtracts the log of each item's sampling probability from its logit.
import torch
import torch.nn.functional as F
def in_batch_softmax_loss(user_vecs, item_vecs, item_log_q, temperature=0.05):
# user_vecs, item_vecs: [B, d], row i is a logged positive pair.
# item_log_q: [B], log of each item's estimated sampling probability in the batch stream.
u = F.normalize(user_vecs, dim=-1)
v = F.normalize(item_vecs, dim=-1)
logits = u @ v.T / temperature # [B, B]: user i scored against every item in the batch
logits = logits - item_log_q[None, :] # sampling-bias correction for popular items
labels = torch.arange(u.shape[0], device=u.device) # the diagonal holds the positives
return F.cross_entropy(logits, labels)Two details matter in practice. If the same item appears twice in a batch, mask the duplicate off-diagonal entry so the model is not taught that an item is a negative for itself. And estimate item_log_q from streaming frequency counts rather than the full catalogue distribution, because what matters is how often an item shows up in training batches.
Serving retrieval
At serving time the item tower runs in a batch job, writes vectors with a model version, and the index is rebuilt or incrementally updated. The user tower runs inside the request. The two must come from the same training run: a user vector from model version 12 queried against item vectors from version 11 produces scores from two unrelated vector spaces, and nothing raises an error. Store the model version with the index and refuse queries that do not match.
def retrieve(user_features, k=500):
model = registry.current("two_tower") # holds version and user tower
index = ann_indexes.get(model.version) # index built from the same version
if index is None:
raise RetrievalUnavailable(model.version) # caller falls back to other sources
u = model.user_tower(user_features) # one forward pass per request
ids, scores = index.search(u, k=k, ef_search=128) # approximate top-k by inner product
return [Candidate(i, s, source="two_tower") for i, s in zip(ids, scores)]Index freshness is the other operational concern. New items have no vector until the next item-tower run, so a catalogue with fast turnover needs either frequent incremental inserts or a separate fresh-content source that does not depend on embeddings. The arithmetic behind vector index memory and recall is in FAISS math.
The ranker: predicting several outcomes
The ranker sees perhaps a thousand candidates and can afford a bigger model with cross features, meaning features that describe the user and the item together: how often this user watched this creator, the time since the user last saw this category, the item's click rate in the user's region. It is usually a multi-task network that predicts several outcomes per item, such as probability of click, expected watch time, probability of a like and probability of a hide or report.
Those predictions are combined into one score with a value function that encodes what the product wants. The weights are a product decision, tuned by online experiments rather than by offline loss.
def value(pred, w):
# pred: model outputs for one (user, item) pair; w: weights chosen by experiment.
return (w.click * pred.p_click
+ w.watch * pred.p_click * pred.expected_watch_seconds / 60.0
+ w.like * pred.p_like
- w.hide * pred.p_hide) # negative outcomes subtract
ranked = sorted(candidates, key=lambda c: value(ranker.predict(user, c), weights), reverse=True)Notice that watch time is conditioned on the click: the model predicts how long the user watches given that they open the item, and the score multiplies that by the probability that they open it. Mixing unconditional and conditional predictions is a common source of scores that look sensible offline and behave strangely in production.
Re-ranking: the list, not the item
Ranking scores items one at a time. A good page is a property of the whole list: ten near-identical videos from one creator may each score highly and still make a poor feed. Re-ranking takes the top few hundred and builds the final list with constraints that need the whole page in view: diversity across creators and topics, a cap on repeats, freshness slots, ads or promoted positions, and policy rules. A simple and effective diversity method is maximal marginal relevance, which trades each item's score against its similarity to items already chosen.
def mmr(items, score, sim, k=20, lam=0.7, max_per_creator=2):
chosen, per_creator = [], {}
pool = list(items)
while pool and len(chosen) < k:
def gain(it):
redundancy = max((sim(it, ch) for ch in chosen), default=0.0)
return lam * score[it.id] - (1 - lam) * redundancy
best = max(pool, key=gain)
pool.remove(best)
if per_creator.get(best.creator, 0) >= max_per_creator:
continue # hard rule beats score
chosen.append(best)
per_creator[best.creator] = per_creator.get(best.creator, 0) + 1
return chosenKeep the re-ranker's rules in configuration with an owner for each, and log which rule moved which item.
The data loop: log what you served
Every model in the funnel learns from impressions joined with outcomes. The most important design decision in the data pipeline is to log the features the ranker actually saw at serving time, keyed by a request id, rather than recomputing them later from the feature store. Recomputed features differ from served ones because counters have moved on, so the model trains on values it will never see in production. This training-serving skew is silent: offline metrics stay good while online behaviour drifts.
- Impression log. Request id, user id, item id, position, source, model versions and the full feature vector as served.
- Event stream. Clicks, watch time, likes, hides, each with the request id and item id of the impression that produced it.
- Join. A windowed join attributes events to impressions. Pick the window per outcome; watch time needs longer than a click.
- Labels. Unclicked impressions become negatives only after the window closes, otherwise late clicks are mislabelled.
The loop also creates a feedback problem. The system only learns about items it showed, and it showed what the previous model liked. Without deliberate exploration, a small slice of traffic spent on items the model is unsure about, the system narrows over time and new items never collect enough data to compete. Multi-armed bandits give a principled way to spend that exploration budget.
Worked example: one request
A user opens a video app at 21:00. The request carries the user id, locale and device. The service fetches about 200 user features from the online feature store in 8 milliseconds. Candidate sources run in parallel: the user tower runs and the ANN index returns 800 items in 6 milliseconds; co-visitation on the last five watched videos returns 600; follows return 150 recent uploads; trending for the locale returns 300; exploration samples 100. After merge, deduplication and filters for already-seen and region-blocked items, 1,100 candidates remain.
The ranker fetches item and cross features in batch and scores the 1,100 candidates in 35 milliseconds on CPU. The value function combines its five heads. The re-ranker takes the top 300, applies MMR with at most two items per creator, reserves one slot for fresh content and one for an exploration item, and returns 20. The whole path takes about 80 milliseconds. Every one of the 20 impressions is logged with its features and the request id; the user watches three videos and hides one, and those events join back to the impressions within the hour for the next training run. The figures are illustrative, chosen to show proportions.
Evaluation and experiments
Offline metrics tell you whether a model is broken, rarely whether it is better. Measure retrieval with recall at k on held-out next interactions and the ranker with log loss per head, then ship through an online experiment with guardrails such as hides and session length. See A/B testing.
Failure modes
- Version skew between towers. User and item vectors from different training runs produce nonsense rankings with no error. Pin the index to the model version.
- Training-serving skew. Features recomputed for training differ from those served. Log features at serving time.
- Popularity collapse. Without bias correction and exploration, the same popular items fill every feed.
- Cold start. New users and items have no history. Use content features in the towers, a fresh-content source and onboarding signals.
- Source outage. If the ANN index is down, the request should degrade to the other sources, not fail. Test that path.
- Rule sprawl. Re-ranking rules accumulate without owners and quietly cancel each other out.
Trade-offs
| Decision | Option A | Option B |
|---|---|---|
| Retrieval | One two-tower source: simple, but misses recency and social signals | Several sources with quotas: better recall, more pipelines to operate |
| Ranker size | Bigger model on fewer candidates: better precision, recall capped upstream | Smaller model on more candidates: more recall, weaker ordering |
| Freshness | Daily retraining: cheap, stale for fast catalogues | Streaming updates: fresh, harder to validate and roll back |
| Objective | Single click target: easy to tune, rewards clickbait | Multi-task value function: closer to product goals, weights need experiments |
What to do next
- Write down the latency budget and catalogue size, and work out how many items each stage can afford to score.
- Add at least two candidate sources beyond embedding retrieval, each with a quota, and measure recall at k for each.
- Train a two-tower model with in-batch softmax and log-frequency correction, and pin the ANN index to the model version.
- Log features as served, keyed by request id, and build training data by joining events to impressions.
- Define the ranker's value function explicitly and tune its weights with online experiments and guardrails.
- Reserve a small, measured slice of traffic for exploration, and test the degraded path when a source fails.