System Design Cases
Recommendation System
Recommendation system case study (Netflix/YouTube/Spotify class). Two-stage funnel: candidate generation via two-tower ANN over millions of items, then ranking via DLRM on top 1000, then re-rank for diversity and business rules. Includes cold start (popular + demo cohort + bandit explore), real-time signal updates via Kafka+Flink streaming, A/B testing with experiment assigner, drift detection and retrain pipeline. 5 scenarios and 2 ADRs (two-stage vs single model, real-time vs batch features).
Рекомендательная система
Схема показывает не «магическую ML-модель», а две связанные системы: online serving выдаёт versioned slate, а evidence pipeline фиксирует, что пользователь действительно увидел и сделал. Без второго контура нельзя честно обучать модель, измерять experiment или объяснять регрессию.
Online request
- API gateway аутентифицирует запрос и передаёт его orchestrator.
- Experiment service стабильно назначает variant по заранее объявленной unit.
- Feature service возвращает значения вместе с event time и версией.
- Candidate service делает approximate retrieval из совместимого ANN index.
- Ranker оценивает кандидатов.
- Policy применяет eligibility, safety, inventory и diversity ограничения.
- Перед ответом server логирует slate: request, версии, candidates, scores, positions и selection propensity.
Candidate count, ANN algorithm, efSearch, latency budget и recall —
benchmark-параметры. HNSW не обещает универсальные 25 ms, 97% recall или
строгую сложность O(log N) для любого распределения, filter workload,
hardware и шардирования.
Impression — не response
Ответ сервера лишь предлагает slate. Exposure появляется, когда client реально отрендерил item согласно продуктовой спецификации видимости. Client отправляет его через authenticated telemetry API и gateway; прямого доступа клиента к Kafka нет.
Раздельно хранятся:
- assignment и trigger experiment;
- server slate и позиции;
- client render exposure;
- action (
click,watch,dismiss, conversion) с event time; - propensity/selection probability для контролируемого exploration;
- версии model, index, feature schema, policy и client.
Так можно находить telemetry loss, position/presentation bias и ошибочные joins. Implicit feedback не является unbiased label: пользователь не может кликнуть то, что ему не показали.
Streaming features без потерянных обновлений
Feature job читает Kafka как consumer, применяет event-time dedupe и watermark,
а затем пишет versioned value через compare-and-set/idempotent key. Нельзя делать
неусловный GET; base + delta; SET: параллельные updates потеряют изменения.
Session signal и долговременный user embedding лучше хранить разными features.
Формула base_embedding + delta допустима только как versioned model contract,
обученный именно для такой композиции; это не общее свойство embeddings.
Training-serving consistency
Feature store помогает, но не устраняет skew автоматически. Training dataset:
- выполняет point-in-time join по времени каждого example;
- не видит future/backfilled values, которые тогда ещё не были доступны;
- использует те же definitions или сравниваемую shared representation;
- содержит logged serving features на sample для прямой проверки skew;
- связывает slate, exposure и action по стабильным ids и допустимому окну;
- учитывает position/exposure bias и sampling probability;
- фиксирует data/code/schema/objective versions.
Offline split следует времени: validation/test идут позже training window. Leakage или mismatched feature semantics делают красивую offline metric непригодной для release.
Совместный release модели и индекса
Query/user encoder, item embeddings, ANN index, ranker, feature schema и policy образуют совместимый bundle. Нельзя отдельно раскатить новый encoder поверх старого index и надеяться на semantic compatibility.
Release coordinator:
- проверяет offline quality, calibration, slices и policy constraints;
- строит index и измеряет recall/latency/memory на production-like corpus;
- подписывает manifest со всеми version IDs и checksums;
- делает shadow или canary;
- проверяет telemetry, latency, errors и guardrails;
- постепенно ramp-ит либо атомарно pin-ит прошлый bundle при rollback.
Incremental index update может ухудшать качество или фрагментировать graph; политика rebuild определяется измерениями, а не обещанием HNSW.
Cold start
Новый пользователь не требует угадывать чувствительные demographic attributes. Fallback может смешивать:
- явно выбранные topics;
- текущий non-sensitive context;
- свежие eligible popular items;
- content similarity без user history;
- небольшую долю контролируемого exploration с logged propensity.
Смесь вроде 70/20/10 не является универсальным рецептом. Её параметры,
privacy basis и success metrics проходят experiment. Тот же fallback покрывает
timeout feature store или пустой candidate set и всегда применяет eligibility.
Иллюстративная ёмкость
Чтобы числа можно было проверить, зададим assumptions:
600 млн DAU;5sessions в день;4recommendation requests на session;20items в slate;- в среднем
25%items действительно становятся exposures; 6 млрдaction events в день как отдельная заданная оценка;- peak =
3×среднего.
Тогда:
- requests/day:
600M × 5 × 4 = 12 млрд; - average request rate:
12 млрд / 86 400 ≈ 138 889 rps; - design peak:
≈ 416 667 rps; - returned placements:
12 млрд × 20 = 240 млрд/день; - measured exposures при assumption 25%:
60 млрд/день; - average action event rate:
6 млрд / 86 400 ≈ 69 444/s.
Returned placement нельзя называть impression без client render signal.
Для 5 млрд items и 256 float32-компонент raw vector занимает
256 × 4 = 1024 B; только vectors — 5,12 TB в десятичном измерении.
Graph links, ids, metadata, filters, replicas и alignment добавляются отдельно.
Фраза «индекс равен 30 TB» может быть capacity estimate конкретной конфигурации,
но не выводом из размерности.
Эксперименты
До чтения effect experiment проверяет:
- stable assignment и declared randomization unit;
- sample-ratio mismatch (SRM);
- telemetry completeness и invariant metrics;
- contamination, carry-over и mutually exclusive layers;
- заранее объявленные primary metric, guardrails, MDE и horizon;
- fixed-horizon или корректную sequential policy, если результаты смотрят рано.
За 48 часов можно принять решение о коротких safety/latency guardrails, но нельзя
заявить улучшение day-7 retention: окно ещё не закрылось. Ramp и долгосрочный
causal readout — разные решения. Если SRM найден, effect не интерпретируется до
устранения причины.
Drift и инциденты качества
Падение CTR — симптом, а не доказательство model drift. Возможные причины: изменившийся traffic mix, сезонность, сломанная exposure telemetry, UI release, feature staleness, index mismatch, experiment interaction или сама модель.
Monitor сначала сопоставляет data/model/index/feature/client versions и slices. По evidence команда может:
- исправить telemetry и не трогать model;
- rollback-нуть несовместимый bundle;
- обновить index;
- retrain-ить на корректном point-in-time dataset;
- признать ожидаемую сезонность и изменить alert baseline.
Автоматический «CTR упал — retrain на последних 7 днях — deploy» способен закрепить corrupted labels и feedback loop, поэтому release проходит evidence gate.
Failure paths
| Сбой | Поведение |
|---|---|
| Feature store timeout | bounded deadline, eligible fallback, отдельная метрика fallback rate |
| ANN shard unavailable | partial candidate set только при quality floor; иначе fallback |
| Ranker timeout | deterministic lightweight ranking над кандидатами |
| Telemetry Kafka недоступна | collector durable-buffer/retry; не выдаёт потерянную запись за принятую |
| Duplicate action | idempotent event id и event-time dedupe |
| Late action | попадает в допустимое label window либо явно исключается с disposition |
| Bad model/index bundle | canary stops; pin previous complete manifest |
| SRM | causal scorecard blocked until diagnosis |
Privacy требует purpose limitation, retention, deletion propagation, access control и исключения запрещённых sensitive/proxy features. Controlled exploration не отменяет safety policy.
Сценарии
Двухстадийный online serving возвращает slate и логирует совместимые версии.
Client сообщает фактический render через API; feature и training pipelines читают один durable event независимо.
Новый пользователь или dependency timeout получает privacy-aware fallback.
SRM блокирует интерпретацию, а long-horizon metric ждёт полного окна.
Point-in-time dataset, model, ANN index и manifest проходят единый release gate.
Monitor исследует telemetry, traffic, skew и versions до rollback или retrain.
Проверяемые метрики
- p50/p95/p99 latency отдельно для features, retrieval, ranking и policy;
- ANN recall@K на representative queries, filtered recall и empty-set rate;
- feature freshness, missingness и online/offline skew;
- slate-to-render ratio, event loss, dedupe и late-event rate;
- experiment SRM, guardrails, power и закрытие declared horizon;
- outcome slices, calibration, diversity и safety violations;
- bundle compatibility, canary rollback time и fallback rate.
SLO задаются продуктом и нагрузочными тестами; архитектура не обещает конкретное качество или причинный uplift.
Первичные источники
- Google Research, two-stage YouTube recommendations: https://research.google/pubs/deep-neural-networks-for-youtube-recommendations/
- Malkov, Yashunin, HNSW: https://arxiv.org/abs/1603.09320
- Feast, point-in-time joins: https://github.com/feast-dev/feast/blob/master/docs/getting-started/concepts/point-in-time-joins.md
- Google, Rules of Machine Learning: https://developers.google.com/machine-learning/guides/rules-of-ml/
- Google Research, propensity and implicit-feedback bias: https://research.google/pubs/attribute-based-propensity-for-unbiased-learning-in-recommender-systems-algorithm-and-case-studies/
- Microsoft Research, SRM diagnosis: https://www.microsoft.com/en-us/research/articles/diagnosing-sample-ratio-mismatch-in-a-b-testing/