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 или объяснять регрессию.
Candidate count, ANN algorithm, efSearch, latency budget и recall —
benchmark-параметры. HNSW не обещает универсальные 25 ms, 97% recall или
строгую сложность O(log N) для любого распределения, filter workload,
hardware и шардирования.
Ответ сервера лишь предлагает slate. Exposure появляется, когда client реально отрендерил item согласно продуктовой спецификации видимости. Client отправляет его через authenticated telemetry API и gateway; прямого доступа клиента к Kafka нет.
Раздельно хранятся:
click, watch, dismiss, conversion) с event time;Так можно находить telemetry loss, position/presentation bias и ошибочные joins. Implicit feedback не является unbiased label: пользователь не может кликнуть то, что ему не показали.
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.
Feature store помогает, но не устраняет skew автоматически. Training dataset:
Offline split следует времени: validation/test идут позже training window. Leakage или mismatched feature semantics делают красивую offline metric непригодной для release.
Query/user encoder, item embeddings, ANN index, ranker, feature schema и policy образуют совместимый bundle. Нельзя отдельно раскатить новый encoder поверх старого index и надеяться на semantic compatibility.
Release coordinator:
Incremental index update может ухудшать качество или фрагментировать graph; политика rebuild определяется измерениями, а не обещанием HNSW.
Новый пользователь не требует угадывать чувствительные demographic attributes. Fallback может смешивать:
Смесь вроде 70/20/10 не является универсальным рецептом. Её параметры,
privacy basis и success metrics проходят experiment. Тот же fallback покрывает
timeout feature store или пустой candidate set и всегда применяет eligibility.
Чтобы числа можно было проверить, зададим assumptions:
600 млн DAU;5 sessions в день;4 recommendation requests на session;20 items в slate;25% items действительно становятся exposures;6 млрд action events в день как отдельная заданная оценка;3× среднего.Тогда:
600M × 5 × 4 = 12 млрд;12 млрд / 86 400 ≈ 138 889 rps;≈ 416 667 rps;12 млрд × 20 = 240 млрд/день;60 млрд/день;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 проверяет:
За 48 часов можно принять решение о коротких safety/latency guardrails, но нельзя
заявить улучшение day-7 retention: окно ещё не закрылось. Ramp и долгосрочный
causal readout — разные решения. Если SRM найден, effect не интерпретируется до
устранения причины.
Падение CTR — симптом, а не доказательство model drift. Возможные причины: изменившийся traffic mix, сезонность, сломанная exposure telemetry, UI release, feature staleness, index mismatch, experiment interaction или сама модель.
Monitor сначала сопоставляет data/model/index/feature/client versions и slices. По evidence команда может:
Автоматический «CTR упал — retrain на последних 7 днях — deploy» способен закрепить corrupted labels и feedback loop, поэтому release проходит evidence gate.
| Сбой | Поведение |
|---|---|
| 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.
SLO задаются продуктом и нагрузочными тестами; архитектура не обещает конкретное качество или причинный uplift.
Введите числа или выберите пресет