RecSys · учебник
Тренажёр Виджеты Повторение О проекте Все главы ← Обзор Retrieval →

Дополнение · X-алгоритм · страница 2 из 11

Пайплайн: как лента собирается из типовых стадий

Лента «Для вас» описана как последовательность из 75 стадий шести типов. Разбираем фреймворк, на котором это держится: что выполняется параллельно, а что обязано идти по порядку, что происходит при падении источника, и почему гидратация запроса разделена на два круга.

Коротко

  • Шесть типов стадий: гидратор запроса, источник, гидратор кандидатов, фильтр, скорер, селектор и побочный эффект. Каждая стадия — структура с тремя методами, и всё.
  • У каждой стадии есть выключатель — метод enable(query), решающий по конкретному запросу, выполняться ли. Через него работают эксперименты.
  • Источники и гидраторы бегут параллельно, фильтры и скореры — последовательно. И это не случайность: у первых нет зависимостей, у вторых есть.
  • Падение источника не роняет запрос. Ошибка молча отбрасывается, лента собирается из того, что успело вернуться.
  • Побочные эффекты запускаются после отправки ответа и не задерживают пользователя.

1. Шесть типов стадий

Весь фреймворк — это 1234 строки на семь файлов. Каждый файл описывает один тип стадии, и все они устроены одинаково.

pub trait Filter<Q, C>: Any + Send + Sync {
    fn enable(&self, _query: &Q) -> bool { true }
    fn filter(&self, query: &Q, candidates: Vec<C>) -> FilterResult<C>;
    fn name(&self) -> &'static str { ... }
}

Фрагмент из filter.rs · код X, Apache 2.0, коммит 28e414f

Три метода: включён ли, что делает, как называется. Имя нужно для метрик — по нему считается, сколько эта стадия отбросила и сколько заняла. Метод run, который вызывает пайплайн, — это обёртка вокруг filter, добавляющая замеры.

Тип стадииЧто делаетПример из ленты
Гидратор запросаДотягивает данные к запросу до того, как появились кандидатыСписок подписок, блокировок, уже показанных постов
ИсточникВозвращает кандидатовСвежие посты подписок, retrieval, кластерная похожесть
Гидратор кандидатовДотягивает данные к каждому кандидатуТекст поста, автор, счётчики, язык, semantic ID
ФильтрВыбрасывает часть кандидатовСтарше 48 часов, свои посты, уже показанные
СкорерПроставляет числаМодель, сведение в скор, переранжирование
СелекторОтбирает и упорядочиваетВзять топ-K, смешать с рекламой
Побочный эффектДелает что-то после ответаЗаписать показы, обновить кеш, отправить события
Почему это важнее, чем кажется

Такая типизация выглядит рутинной инженерией, но именно она делает возможным всё остальное в системе.

Добавить источник — значит написать одну структуру. Не трогая ничего, что уже работает. Тридцать фильтров в ленте существуют потому, что добавление тридцать первого стоит десяток строк.

Замеры одинаковы для всех. Обёртка run считает время и размеры для любой стадии, поэтому в системе автоматически есть ответ на вопрос «где мы теряем время» и «какой фильтр сколько выбрасывает» — без единой строки специального кода.

Эксперименты ставятся на уровне стадии. Метод enable получает запрос и решает по нему; внутри он читает параметры. Значит, любую стадию можно включить на процент трафика, не разветвляя код.

Сравните с тем, как мы описывали архитектуру рантайма: там шла речь о том, что рекомендательный сервис — это конвейер, и что его сила в единообразии стадий. Здесь это доведено до библиотеки.

2. Порядок выполнения

Главная функция пайплайна читается как оглавление всей ленты:

let hydrated_query = self.hydrate_query(query).await;
let hydrated_query = self.hydrate_dependent_query(hydrated_query).await;
let candidates = self.fetch_candidates(&hydrated_query).await;
let hydrated_candidates = self.hydrate(&hydrated_query, candidates).await;
let (kept, mut filtered) = self.filter(&hydrated_query, hydrated_candidates.clone());
let scored = self.score(&hydrated_query, kept).await;
let SelectResult { selected, non_selected } = self.select(&hydrated_query, scored);
let post_selection = self.hydrate_post_selection(&hydrated_query, selected).await;
let (mut final_candidates, post_filtered) = self.filter_post_selection(&hydrated_query, post_selection);

Фрагмент из candidate_pipeline.rs · код X, Apache 2.0, коммит 28e414f

1. Гидратация запроса два круга, внутри круга — параллельно 2. Источники все сразу, параллельно; упавшие просто молчат 3. Гидратация кандидатов 12 гидраторов параллельно: текст, автор, счётчики… 4. Фильтры до скоринга строго по очереди — каждый видит выход прошлого 5. Скореры по очереди: модель → сведение → переранжирование 6. Селектор сортировка и обрезка до топ-K 7. Гидратация после отбора дорогое — только для тех, кто прошёл 8. Фильтры после отбора видимость, дедуп разговоров ответ пользователю и только потом — побочные эффекты Синим и оранжевым — то, что выполняется параллельно. Розовым и фиолетовым — то, что обязано идти строго по очереди. Дорогая гидратация стоит после отбора: спрашивать про сотню отобранных дешевле, чем про тысячу кандидатов.
Восемь шагов и 75 стадий внутри них. Порядок задан не вкусом, а зависимостями между данными.

Что параллельно, а что нет

Разница видна прямо в коде и она не косметическая.

Источники — параллельно:

let source_futures = sources.iter().map(|s| s.run(query));
let results = join_all(source_futures).await;

Фильтры — последовательно:

for filter in enabled {
    let result = filter.run(query, candidates);
    ...
}

Фрагмент из candidate_pipeline.rs · код X, Apache 2.0, коммит 28e414f

Почему так? Источники друг от друга не зависят: каждый берёт запрос и возвращает список. Их можно запустить одновременно, и общее время станет временем самого медленного, а не суммой.

Фильтры зависят: каждый получает то, что осталось после предыдущего. Запустить их параллельно на одном и том же входе можно было бы — результат в теории тот же, — но тогда каждый фильтр обрабатывал бы полный список вместо усечённого. При тридцати фильтрах и тысячах кандидатов это дороже, а не дешевле. Плюс порядок несёт смысл: дешёвые фильтры стоят раньше и уменьшают работу для дорогих.

Арифметика параллельности

Пусть источников семь, и они отвечают за 3, 5, 8, 4, 12, 6 и 2 миллисекунды. Последовательно это 40 мс, параллельно — 12 мс, то есть время самого медленного.

Отсюда важное следствие для проектирования: оптимизировать надо не средний источник, а худший. Ускорив тот, что отвечает за 3 мс, вдвое, вы не выиграете ничего. Это стандартный аргумент про хвостовые задержки из главы «Архитектура рантайма», и здесь он записан прямо в структуре кода.

3. Что происходит, когда источник падает

Строка, которую легко пропустить, но которая определяет поведение системы под нагрузкой:

let results = join_all(source_futures).await;

let mut collected = Vec::new();
for mut candidates in results.into_iter().flatten() {
    collected.append(&mut candidates);
}

Фрагмент из candidate_pipeline.rs · код X, Apache 2.0, коммит 28e414f

Каждый источник возвращает Result<Vec<C>, String> — либо кандидатов, либо ошибку. А flatten по итератору результатов молча выбрасывает все ошибки и оставляет только успешные.

Значит, если сервис свежих постов подписок упал, запрос не падает. Лента соберётся из того, что вернули остальные источники: она станет хуже — не будет постов от тех, на кого вы подписаны, — но она будет.

Правильное ли это решение

Да, и это стандартный подход к отказоустойчивости в рекомендательных сервисах, который мы разбирали в главе «Отказоустойчивость» под названием «лестница фоллбеков». Рекомендация — не банковская транзакция: показать неидеальную ленту лучше, чем не показать ничего.

Но у решения есть цена, и её стоит назвать вслух. Деградация тихая. Пользователь не узнает, что половина источников молчала, и по одному запросу это не отличить от нормальной работы. Если бы источник отваливался в двух процентах запросов, никакой алерт бы не сработал — просто у части людей лента была бы систематически хуже.

Именно поэтому в обёртке run каждой стадии стоят замеры: спасает не сам факт устойчивости, а то, что рядом считается, сколько раз стадия отработала и сколько вернула. Устойчивость без наблюдаемости превращается в тихую поломку.

4. Гидратация запроса: два круга

Обратите внимание, что гидратация запроса вызывается дважды:

let hydrated_query = self.hydrate_query(query).await;
let hydrated_query = self.hydrate_dependent_query(hydrated_query).await;

Фрагмент из candidate_pipeline.rs · код X, Apache 2.0, коммит 28e414f

Первый круг — гидраторы, которым нужен только исходный запрос: список подписок, блокировок, скрытых аккаунтов, демографические данные. Все они параллельны друг другу.

Второй круг — те, которым нужны результаты первого. Например, чтобы посчитать взаимные подписки, надо сначала получить список подписок. Или чтобы построить последовательность действий для модели, надо знать, какие посты уже показывались.

Это простейшая форма планирования зависимостей: вместо полноценного графа — два уровня. Решение прагматичное: настоящий граф зависимостей потребовал бы описывать связи между стадиями, а два уровня покрывают почти все реальные случаи и не требуют ничего, кроме размещения гидратора в нужный список.

Что именно дотягивается к запросу

Гидраторов запроса двадцать. Вот наиболее показательные.

ГидраторЧто приноситЗачем
ScoringSequenceQueryHydratorПоследовательность действий для ранжирующей моделиГлавный вход модели — та самая история из x04
RetrievalSequenceQueryHydratorПоследовательность для retrievalОтдельная, потому что модель другая
FollowedUserIdsQueryHydratorСписок подписокНужен источнику постов подписок и разметке in-network
BlockedUserIdsQueryHydrator и Muted…Блокировки и скрытые аккаунтыФильтрация
MutualFollowQueryHydratorВзаимные подпискиПрибавка к весу ответа: разговор со знакомым ценится отдельно
ImpressionBloomFilterQueryHydratorФильтр Блума с показанными постамиНе показывать дважды
UserDemographicsQueryHydratorСтрана, язык, возрастПризнаки модели
UserInferredGenderQueryHydratorПредполагаемый полПризнак модели
Фильтр Блума — прямо на дорожке запроса

Строка ImpressionBloomFilterQueryHydrator заслуживает отдельной остановки, потому что это тот самый фильтр Блума из главы «Фильтр Блума и пагинация», встреченный в проде.

Задача: не показывать пост, который человек уже видел. Наивное решение — таскать за пользователем полный список показанных идентификаторов. При активном чтении это десятки тысяч чисел на каждый запрос — по объёму данных и по времени передачи неприемлемо.

Фильтр Блума даёт компактную структуру на несколько килобайт с односторонней ошибкой: если он говорит «не видел» — точно не видел, если «видел» — с малой вероятностью ошибается. Ошибка стоит нам того, что мы лишний раз спрячем хороший пост. Ошибка в обратную сторону — показ дубликата — заметна пользователю куда сильнее.

Посчитать, во что обходится такой размен при конкретных размерах, можно в виджете фильтра Блума. Любопытно, что в ленте рядом с фильтром Блума работают ещё три фильтра по уже показанному — про это на странице фильтров.

5. Две трубы: посты и всё остальное

Пайплайнов на самом деле два, и они вложены друг в друга.

Post Pipeline — то, что мы разбирали выше: 75 стадий, находящих и упорядочивающих посты.

Blending Pipeline — обёртка вокруг него. Для неё отранжированные посты — всего лишь один источник среди прочих. Другие источники приносят то, что моделью не ранжируется: рекламу, блок «кого читать», промо-подсказки.

Blending Pipeline — собирает то, что человек увидит Post Pipeline 75 стадий источники → гидратация → фильтры → скоринг → отбор → фильтры = отранжированные посты реклама кого читать промо и подсказки BlenderSelector чередует по правилам, фиксированные позиции лента на экране
Внешний пайплайн видит внутренний как обычный источник кандидатов. Тот же фреймворк, вложенный сам в себя.
Почему это удачное разделение

Реклама и рекомендации живут по разным законам. У рекламы свой аукцион, свои обязательства перед рекламодателем, свои ограничения по соседству. Пытаться выразить это через тот же скор, что и органические посты, — заведомо провальная затея: у них нет общей единицы измерения.

Разделив трубы, систему избавляют от этой задачи вовсе. Ранжирование постов отвечает на вопрос «какой пост лучше какого». Блендинг отвечает на другой — «на какие позиции ставить не-посты». Второй вопрос решается правилами, а не моделью, и это правильно: доля рекламы в ленте — решение бизнеса, а не предсказание.

Мы приходили к тому же выводу в главе «Как проектировать систему с нуля»: слои с разными целевыми функциями не сводят в один скор, их разносят по стадиям.

6. Побочные эффекты: после ответа

fn run_side_effects(&self, input: Arc<SideEffectInput<Q, C>>) {
    tokio::spawn(... async move {
        let futures = side_effects.iter()
            .filter(|se| se.enable(input.query.clone()))
            .map(|se| se.run(input.clone()));
        let _ = join_all(futures).await;
    });
}

Фрагмент из candidate_pipeline.rs · код X, Apache 2.0, коммит 28e414f

Ключевое здесь — tokio::spawn и let _ =. Задача запускается в фоне, её результат не ожидается и не проверяется. Ответ пользователю уже отправлен.

Что делают побочные эффекты в ленте: записывают, какие посты были показаны (чтобы не показать снова), обновляют кеш отранжированных постов, отправляют события про рекламу, пишут статистику по авторам и логи экспериментов.

Все они обязаны быть вне критического пути. Запись показов может занять десятки миллисекунд — заставлять пользователя ждать этого ради того, что нужно только следующему запросу, бессмысленно.

Обратная сторона

Побочные эффекты — это место, где рождается рассогласование между тем, что показано, и тем, что залогировано. Ответ ушёл, а запись показов ещё не произошла или произошла с ошибкой, которую никто не проверил: let _ = прямо говорит, что результат игнорируется.

Практическое следствие видно в самой ленте: там два независимых фильтра по уже показанным постам, читающих из разных источников, плюс фильтр Блума в гидраторах. Три механизма на одну задачу — это не паранойя разработчиков, а прямое признание того, что запись показов ненадёжна.

Частые ошибки и подводные камни

На чём спотыкаются
  • Считать, что фильтры можно распараллелить. Формально да, но тогда каждый обрабатывает полный список вместо усечённого предыдущими. При тридцати фильтрах это дороже.
  • Не замечать, что ошибки источников проглатываются. flatten выбрасывает Err молча. Система устойчива, но деградирует тихо — без метрик по каждой стадии это неотличимо от нормы.
  • Думать, что дорогая гидратация идёт до фильтров. Наоборот: дешёвое и массовое — до, дорогое и точечное — после отбора топ-K.
  • Считать блендинг частью ранжирования. Это отдельная труба, для которой все отранжированные посты — один источник среди прочих.
  • Полагаться на побочные эффекты как на гарантию. Они запускаются после ответа, их результат не проверяется. Отсюда три параллельных механизма защиты от повторного показа.

Вопросы с собеседований

Спроектируйте пайплайн рекомендательного сервиса. Какие стадии и в каком порядке?

Гидратация запроса → источники кандидатов → гидратация кандидатов → дешёвые фильтры → скоринг → отбор топ-K → дорогая гидратация → дорогие фильтры → ответ → побочные эффекты в фоне.

Ключевые обоснования порядка: источники не зависят друг от друга и идут параллельно; фильтры зависят и идут по очереди, от дешёвых к дорогим; всё, что требует похода в другой сервис по паре «объект и пользователь», ставится после отбора, потому что там на порядок меньше объектов; всё, что нужно не текущему запросу, а следующему, уходит в фон после отправки ответа.

Что должно происходить, если один из источников кандидатов недоступен?

Запрос должен выполниться на оставшихся источниках. Рекомендация — не транзакция: неполная лента лучше ошибки. В разобранном коде это буквально одна строка: результаты источников проходят через flatten, который отбрасывает ошибки.

Но обязательное дополнение: деградация должна быть видимой. Нужны метрики по каждому источнику — сколько раз вызвали, сколько раз получили ошибку, сколько кандидатов вернулось. Без этого падение источника на нескольких процентах трафика не обнаружится никогда, а часть пользователей будет систематически получать худшую ленту.

Зачем нужен метод «включён ли» у каждой стадии?

Для экспериментов. Метод получает запрос и решает по нему — внутри он читает параметры конфигурации, а те раскатываются на процент трафика. Таким образом любую стадию можно включить или выключить на части пользователей, не разветвляя код и не выкатывая новую версию.

Второе применение — условная логика: источник тем работает только для тематических запросов, фильтр видео — только когда клиент попросил исключить видео. Вместо ветвлений внутри стадии условие поднимается на уровень фреймворка, где оно ещё и автоматически попадает в метрики.

Почему гидратация запроса разбита на два круга?

Потому что часть гидраторов зависит от результатов других. Взаимные подписки нельзя посчитать, не получив список подписок; последовательность действий для модели строится с учётом уже показанных постов.

Внутри круга гидраторы параллельны, между кругами — барьер. Это упрощённая замена графу зависимостей: настоящий граф требовал бы описывать связи явно, а два уровня покрывают почти все случаи и не требуют ничего, кроме размещения гидратора в нужный список.

Как измерять, где пайплайн теряет время?

Замеры должны быть встроены во фреймворк, а не расставлены руками по стадиям. В разобранном коде у каждого типа стадии есть обёртка run, которая вызывает содержательный метод и попутно фиксирует время, число входных и выходных кандидатов, а также имя стадии.

Отсюда бесплатно получаются ответы на два главных вопроса эксплуатации: где узкое место по времени и какой фильтр сколько выбрасывает. Причём при параллельном выполнении оптимизировать надо худшую стадию, а не среднюю — общее время равно времени самой медленной.

Шпаргалка одним экраном

Типы стадий

Гидратор запроса, источник, гидратор кандидатов, фильтр, скорер, селектор, побочный эффект.

Интерфейс стадии

Включён ли · что делает · как называется. Обёртка добавляет замеры автоматически.

Параллельно

Источники и гидраторы. Время = время самого медленного.

Последовательно

Фильтры и скореры: каждый видит выход предыдущего.

При падении источника

Ошибка отбрасывается, лента собирается из остального. Деградация тихая — нужны метрики.

Два круга гидратации

Второй для тех, кому нужны результаты первого.

Дорогое — после отбора

Спросить про сотню отобранных дешевле, чем про тысячу кандидатов.

Две трубы

Post Pipeline ранжирует посты; Blending Pipeline подмешивает рекламу и не-посты.

Побочные эффекты

После ответа, в фоне, результат не проверяется. Отсюда дублирование защиты от повторов.

Первоисточники