Как за 5 недель построить рекомендательную систему в TravelTech: Kafka и MongoDB вместо Feature Store. CatBoost.. CatBoost. ClickHouse.. CatBoost. ClickHouse. fastapi.. CatBoost. ClickHouse. fastapi. feature store.. CatBoost. ClickHouse. fastapi. feature store. kafka.. CatBoost. ClickHouse. fastapi. feature store. kafka. mlops.. CatBoost. ClickHouse. fastapi. feature store. kafka. mlops. mongodb.. CatBoost. ClickHouse. fastapi. feature store. kafka. mlops. mongodb. travel tech.. CatBoost. ClickHouse. fastapi. feature store. kafka. mlops. mongodb. travel tech. рекомендательные системы.

Привет! Меня зовут Кристина, я MLOps-инженер в Туту. Занимаюсь тем, что помогаю рекомендательным системам добраться до прода со всеми компромиссами, горящими дедлайнами и новыми идеями. 

Эта статья — про один из таких запусков.

В идеальном ML-мире запуск рекомендаций выглядит примерно так: полгода проектируют хранилище признаков (Feature Store), настраивают, откуда и как берутся данные, гоняют тяжёлые расчёты фичей и моделей, а отдельная команда следит, не деградирует ли модель из‑за изменений в данных. 

В реальном бизнесе у тебя есть 5 недель до старта высокого сезона, два инженера, DS и задача: сделать так, чтобы пользователь, который купил билет, сразу увидел релевантный отель. Рассказываем, как мы собрали работающую RecSys v1 на привычном стеке: Kafka, MongoDB, ClickHouse. При этом мы сознательно отказались от перфекционизма ради скорости.

Инженерный вызов здесь не в масштабе и не в алгоритмах, а в контексте. Cross-sell в travel — это не «похожие товары». 

Пример

Пользователь купил билет Москва → Сочи на 10–17 июля: значит, нужно показать отели именно в Сочи, именно на эти даты. 

Коллаборативная фильтрация без контекста поездки — «похожие пользователи → похожие отели» — не знает ни город, ни даты, ни то, что заказ только что оплачен и его ещё нет в DWH. 

Нужен подход, где контекст конкретной поездки — куда, когда, с кем — задаётся явно до ранжирования.

«Правильный» путь —  Feature Store и полноценная ML-платформа, занял бы 4–6 месяцев. Бизнесу нужно было проверить гипотезу на живом трафике. Мы собрали v1 на том, что уже работало в проде, с некоторыми компромиссами и без иллюзий насчёт идеальной архитектуры. 

Про стек, ограничения и сроки — в следующем разделе.

Ограничения: срок, команда, знакомый стек

Пять недель — это не так мало, если не пытаться построить всё сразу.

К моменту старта в компании была зрелая data-платформа — Kafka с realtime-потоками заказов, ClickHouse с данными по всем вертикалям, поверх которого мы считали агрегаты для модели, Airflow, K8s для деплоя. Всё это уже работало в проде. ML-платформы не было, нужно было добавить model serving — слой, который на живом трафике принимает запрос и возвращает рекомендацию.

Feature Store (Feast) был в плане, но не в проде. Вместо него использовали MongoDB: туда писали всё, что нужно модели на инференсе — заказы из Kafka и batch-фичи из ClickHouse в одном документе на пользователя. Временно, но работает.

Команда маленькая: DS готовил фичи в ClickHouse и модель, параллельно с тем, как мы строили слой данных и serving. Синхронизировались по контракту — что нужно на входе inference, что возвращать на выходе, и работали независимо.

Срок резали по одному критерию: должно быть в проде к фиксированной дате запуска. Если фича не приближала к этому, она уходила в backlog

На старте было три сценария — покупка, брошенная корзина, поиски. В v1 сознательно взяли один: рекомендация на экране подтверждения заказа, сразу после оплаты транспортного заказа. Остальное – «в планах», архитектуру сразу закладывали под расширение — другие точки входа, гостевые сессии, контекст без длинной истории заказов.

В проде на старте работал только сценарий для пользователей с историей: рекомендации видели те, у кого уже была история в Mongo. Гость или пользователь «с нуля» получал пустую выдачу, сервис отвечал успешно, просто без единой рекомендации. Не баг, а сознательный scope: cold start в v1 не включали.

При этом контракт API проектировали не под один сценарий: в запросе уже лежали is_authenticated и state — поездку можно было взять из запроса, если в профиле ещё пусто. 

Session-агрегаты и cold-модель прикрутили позже — тот же Kafka → Mongo → FastAPI, другие только ранкер и фичи.

В первой версии мы обязаны были поднять слой выдачи рекомендаций. Плюс слой данных: заказы из Kafka, фичи из ClickHouse, всё это в Mongo, откуда инференс читает за миллисекунды. И контракт с DS: что на входе, что на выходе. Без этого два человека и один DS просто не сойдутся за пять недель.

Модель для нас на старте — чёрный ящик. ALS, rules, CatBoost — это реализация внутри predict, не предмет статьи. Менялся ранкер — не менялись API, target order и Mongo.

Как мы всё успели: прагматичная архитектура

Что откуда берётся

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

MongoDB — единственная точка чтения при выдаче рекомендаций. Коллекции user_profiles и hotel_catalog живут в одной БД; на запрос API ходит только в Mongo.

Рис. 1. Запись в два темпа (realtime — Kafka, batch — Airflow из ClickHouse), инференс читает только Mongo: user-документ и каталог отелей города. Ranker: ALS в v1 → CatBoost позже (.cbm в S3 — вне data-plane)

Рис. 1. Запись в два темпа (realtime — Kafka, batch — Airflow из ClickHouse), инференс читает только Mongo: user-документ и каталог отелей города. Ranker: ALS в v1 → CatBoost позже (.cbm в S3 — вне data-plane)

Realtime-путь: когда пользователь оплачивает заказ, событие летит в Kafka. Консьюмер забирает его, обрабатывает и пишет в MongoDB в коллекцию user_profiles — к документу конкретного пользователя.

Batch-путь: каждую ночь Airflow-джоба читает ClickHouse и дописывает агрегаты в MongoDB: user-поля в user_profiles, hotel-поля в hotel_catalog. В первом релизе batch был минимальным — в основном ALS-векторы плюс has_orders / order_count у отелей; recency, CTR и большую часть полей нарастили позже под CatBoost, паттерн upsert не менялся.

Online inference: При запросе API читает только Mongo — документ пользователя и каталог отелей города поездки. Без сложных запросов на лету, без Feature Store. Время ответа предсказуемое.

Первый релиз (v1): ранжирование на векторах из Mongo. Скор — скалярное произведение ALS-векторов в памяти; «веса» жили в документах — nightly batch дописывал als_factors в user_profiles и hotel_catalog. Отдельного файла модели и S3 не было. Если векторов не хватало — fallback на top по order_count в городе.

После первого релиза — CatBoost и S3. Когда DS перешёл на CatBoost, ranker стал другим, но слой данных остался тем же: фичи по-прежнему в Mongo, заказы — из Kafka. Добавился .cbm в object storage и загрузка модели один раз при старте пода — не на каждый запрос

Три решения, которые нас спасли

MongoDB вместо Feature Store

Feast мы планировали, но строить его с нуля за 5 недель невозможно. Поэтому MongoDB стала временным online store.

Храним всё, что нужно для выдачи рекомендаций, в одном документе пользователя. Заказы приходят из Kafka в realtime, batch-фичи дописываются ночью. На момент запроса документ содержит и то, и другое.

# model_context_service.py

async def build_context(self, user_id, order_id, ...):

    # Один запрос к Mongo — весь контекст для инференса

    doc = await self.users_storage.get_collection('user_profiles').find_one({"_id": user_id})

    # Находим target order — транспортный заказ, под который подбираем отель

    target_order_raw = self.find_target_order(doc.get('orders', []), order_id)

    # Конвертируем UTC-даты в локальное время пункта прибытия

    target_order = self._apply_local_dates(target_order_raw)

    return BaseContext(user=doc, target_order=target_order, ...)

Один find_one — и у нас есть заказы пользователя и ALS-векторы. Всё лежит в одном документе, который обновляется из двух источников независимо. Срез user-документа в v1 (данные обезличены):

{

  "_id": 1000001,

  "orders": [

    {

      "order_id": "100000002",

      "order_type": "avia",

      "departure_geo_city_id": 2656915,

      "arrival_geo_city_id": 2656874,

      "departure_date": { "$date": "2099-04-12T14:20:00.000Z" },

      "arrival_date": { "$date": "2099-04-12T16:45:00.000Z" },

      "common_status": "paid",

      "passenger_count": 2

    }

  ],

  "als_factors_1": [0.001, -0.0003, 0.0007, 0.0002, -0.0005, ... /* ~50 значений */]

}

orders[] — из Kafka в realtime (target order и контекст поездки). als_factors_1 — ночной batch из ClickHouse; на v1 ими же и ранжировали. Скаляры вроде recency_days и кликов добавили позже, уже под CatBoost; ALS-векторы в Mongo при этом остались как legacy batch-поля.

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

Target order и контракт API

Главное архитектурное решение v1 — не ранкер, а контракт: какую поездку считаем контекстом и что именно отдаём фронту.

Купили билет Москва → Сочи на 10–17 июля — target_order это этот transport-заказ из истории покупок. Из него берем город (locality), даты, гостей. Нет города назначения — отдаём пустой список, ranker не трогаем.

Даты — это отдельная история. Город и заезд/выезд берём из только что купленного билета. Звучит просто, но в системе билет хранится в UTC, а человек живёт в локальном времени пункта прибытия. Самолёт прилетает в Сочи 13 апреля в 01:30 ночи — для пользователя это уже 13-е. В UTC это может быть ещё 12-е вечером. Если взять дату «как в базе», без перевода в местное время, в API уйдёт заезд на 12-е вместо 13-го.

Сервис не упадёт: ответ 200, отели в Сочи, ранкер отработал. Но ночь не та, и это ломает метрики тихо, без алертов.

Поэтому перед инференсом мы переводим дату прилёта в локальное время пункта назначения и уже из неё считаем checkin/checkout для поиска отелей.

Как за 5 недель построить рекомендательную систему в TravelTech: Kafka и MongoDB вместо Feature Store - 2

Фронту не отдаём сырые id — отдаём готовый поиск: город, даты, гости, отели со сроками. Контекст из билета фронт не собирает сам. 

{

  "type": "hotels",

  "score": 0.82341,

  "search_params": [{

    "locality": 2656874,

    "checkin_date": "2099-04-12T00:00:00Z",

    "checkout_date": "2099-04-18T00:00:00Z",

    "number_of_guests": 2,

    "hotel_id_list": [1000042, 1000018, 1000007],

    "score": 0.82341

  }]

}

Два темпа данных: Kafka + ночной ClickHouse

Главная инженерная проблема cross-sell на экране подтверждения заказа: пользователь только что оплатил билет. В хранилище данных его заказ появится только завтра. 

Значит, если читать всё из batch — модель не видит только что купленный билет и не может подобрать под него отель.

Мы разделили данные по темпу изменения:

  • Orders[] — только Kafka. Заказы критичны для свежести: именно по ним определяется target_order. Kafka-консьюмер пишет их в Mongo после оплаты.

  • ALS — ночной batch. На v1 этого хватало для score. Recency, CTR и остальные скаляры добавили позже, когда перешли на CatBoost. Airflow считает их из ClickHouse раз в ночь и дописывает в user_profiles и hotel_catalog.

# base_orders_loader.pyлогика merge при записи realtime-заказа

def mergeorders_with_deduplication(self, existing_orders, new_orders):

    Merge новых заказов с существующими по order_id.

    Побеждает более свежий (по полю ‘created’).

    Один и тот же заказ может прийти несколько раз:

    статус-апдейт, повторная запись из snapshot — всё обрабатывается здесь.

    orders_dict = {o['order_id']: o for o in existing_orders if o.get('order_id')}

    for new_order in new_orders:

        order_id = new_order.get('order_id')

        if not order_id:

            continue

        existing = orders_dict.get(order_id)

        # Перезаписываем только если новый заказ свежее

        if existing is None or self._is_newer_order(new_order.get('created'), existing.get('created')):

            orders_dict[order_id] = new_order

    return list(orders_dict.values())

Kafka — поток событий, а не «одна строка на заказ». Одно и то же бронирование может прийти несколько раз: сначала без статуса оплаты, потом с paid, при повторной доставке сообщения или при заливке истории из snapshot. При каждой записи мы мержим массив по order_id: если заказ уже есть, то оставляем версию с более свежим created, дубликаты не копим. Для cross-sell это важно: в orders[] всегда актуальное состояние поездки, а не три копии одного билета.

Сервис без истории заказов в Mongo для warm path бесполезен: пользователь давно покупает билеты, но пока мы не залили прошлые заказы из snapshot, в orders[] пусто — рекомендации нет, хотя API отвечает 200.

Поэтому до прода мы не рассчитывали, что Kafka сама накопит месяцы жизни. Один раз прогнали snapshot-топики — история заказов примерно за полгода по авиа, ж/д-транспорту, автобусам, отелям — и залили в Mongo батчами по 2000 сообщений, с merge и dedup по order_id. Это заняло порядка 4–6 дней.

Дальше новые оплаты идут через realtime: свежий paid дописывается в тот же документ. Заливка исторических данных — не архив ради архива, а разовая инициализация, без которой warm path в момент запуска для большинства пользователей просто не существовал бы.

Сервис рекомендаций на запросе читает только Mongo: ему не нужно знать, откуда в документ попали заказы (Kafka) и агрегаты (ClickHouse).

Ночной batch: один SQL — одно поле — один upsert

Realtime закрывает свежесть orders[], но профиль пользователя и фичи отелей — это другая задача. Их считает batch-агрегатор: Airflow-джоба раз в ночь гоняет два процесса — UsersAggregator и HotelsAggregator. В первом релизе в них было по несколько SQL (в основном ALS); сейчас ~21 user-поле и ~65 hotel-полей — нарастили под CatBoost, паттерн тот же: один SQL → одно поле → один upsert.

Схема в текущем масштабе:

Как за 5 недель построить рекомендательную систему в TravelTech: Kafka и MongoDB вместо Feature Store - 3

Один агрегат = один .sql + один dataclass. Хотим добавить фичу hotels_recency_days — пишем SQL, регистрируем агрегат, деплоим. Соседние поля в Mongo не трогаем. Для v1 с маленькой командой это было критично: DS и инженеры могли добавлять фичи параллельно, не ломая друг другу документы.

Каждый SQL возвращает пары (user_id | hotel_id, value). Агрегатор читает ClickHouse батчами и делает upsert только своего поля:

# users_aggregator.py — упрощённо

for aggregate in self.aggregates:

    for aggregate_name, batch in self._aggregator.process([aggregate]):

        await self._storage.upsert_aggregate(

            items=batch,

            field_name=aggregate_name,          # например, "hotels_recency_days"

            collection_name='user_profiles',

        )

Под капотом это обычный $set — массив orders[] из Kafka не перезаписывается:

UpdateOne(

    {'_id': user_id},

    {'$set': {'hotels_recency_days': 117}},

    upsert=True,

)

Пример SQL:

-- hotels_recency_days.sql

select

    customer_user_id as user_id,

    dateDiff('day', max(toDate(order_created_at)), toDate(now())) as hotels_recency_days

from hotels_marts.orders_enriched

where is_sold = 1 and customer_user_id > 0

group by user_id

Аналогично считаются и другие поля — recency, цены, клики. В v1 их ещё не было.

На инференсе — два чтения из Mongo: документ пользователя и каталог отелей города. 

Фрагмент hotel-документа в v1 (обезличено):

{

  "_id": 1000001,

  "geo_city_id": 2656915,

  "has_orders": 1,

  "order_count": 15,

  "als_factors_1": [0.002, -0.001, 0.0008, ... /* ~50 значений */]

}

Отказоустойчивость на уровне одного агрегата. Если один nightly-агрегат не отработал, остальные поля обновятся как обычно; на inference пользователь получит последнее записанное значение, а не 500. Сознательный компромиссный подход: доступность важнее свежести каждого поля в каждую ночь.

Результат

За пять недель собрали и задеплоили в прод работающий RecSys v1. Честно говоря, в какой-то момент казалось что не успеем, но вот как это получилось.

Работы шли параллельно: DS считал фичи и SQL в ClickHouse, мы поднимали FastAPI, merge заказов в Mongo и ночные агрегаторы, Kafka дописывала orders[]. Стыковались по JSON-контракту — что на входе инференс, что отдаём фронту. Когда цепочка сошлась, добавили базовый мониторинг.

Пользователь оплатил билет — и на главной, в карточке поездки или во вкладке заказов видит карусель отелей: Сочи, его даты, подобранные под него, а не общий топ. 

Так мы проверяли гипотезу в A/B: контроль — старая подборка, тест — Best Offer с моделью. За месяц эксперимента конверсия из показа в заказ отеля выросла на единицы процентов, GMV – двузначный рост. Фичу раскатили на всех. Бэкенд за пять недель успел довести рекомендацию до экрана, где человек уже в контексте поездки.

Как за 5 недель построить рекомендательную систему в TravelTech: Kafka и MongoDB вместо Feature Store - 4

Сервис отдаёт рекомендации на живом трафике с медианной задержкой ~100ms, p99 < 400ms —для пользователей с накопленной историей заказов, два обращения в Mongo, ранжирование по десяткам отелей города.

Отдельно стоит отметить, что инфраструктура, собранная за несколько недель, потом пережила смену модели без перестройки. 

На старте был ALS: векторы лежали в Mongo, score считался dot product на инференсе — отдельного файла модели и S3 не было. 

Через несколько месяцев DS перешёл на CatBoost: другие фичи, .cbm в object storage, ranker грузится при старте сервиса. Меняли только слой scoring — между чтением Mongo и формированием ответа. 

Kafka по-прежнему писала orders[] в realtime, nightly job — дописывала агрегаты тем же upsert одного поля, API — отдавал те же search_params с городом, датами и списком отелей. Батч успел вырасти с пары ALS-полей до десятков пользовательских и отельных фич ещё до того, как CatBoost попал в прод — слой данных готовился заранее, ranker догонял. Это лучший индикатор того, что слой данных с самого начала был спроектирован правильно — не под конкретную модель, а под задачу.

Рефлексия: что упростили сознательно

Мы не строили Feature Store — и это было правильным решением на старте. MongoDB с одним документом на пользователя дала скорость запуска, но забрала lineage и единый реестр фич. Когда фичей мало и команда маленькая — это терпимо. Когда фич становится больше — отсутствие Feature Store начинает болеть: непонятно, откуда взялась фича, как давно обновлялась, что сломается если её убрать. Для других ML сервисов в компании мы уже внедрили Feast — feature views, lineage, единый реестр фич.

Transport — детерминированные правила, не модель. Это сознательный выбор: rules прозрачны, дебажатся за минуты, не требуют отдельного model lifecycle. Для v1, где нужно было проверить саму концепцию cross-sell, этого достаточно. 

Cold start в продукт не включали — см. выше; для A/B хватило warm-аудитории.

Перед сезоном сделали нагрузочное тестирование: проверили, что инференс с двумя чтениями из Mongo укладывается в SLA.

Самый неприятный баг в RecSys — не 500, а 200 с пустой каруселью. Поэтому на старте observability строили вокруг этого: скорость inference, ошибки и в логах — причина пустой выдачи (rec_empty_reason). 

Это не ML-monitoring в академическом смысле: мы не смотрели, как сдвинулся score или распределение recency_days. Полагались на ночной пересчёт агрегатов — данные обновляются, ranker меняется, а «почему именно этому пользователю пусто» разбирали по логам.

Заключение

Kafka + ClickHouse + MongoDB — не новый стек. Знакомость инструментов сэкономила недели на онбординге; на v1 ушли в слой данных и контракт inference — что именно считаем поездкой и что отдаём фронту.

Два решения сэкономили больше всего времени: MongoDB как временный online store вместо Feature Store и target order как явный контракт API. Первое придётся эволюционировать по мере роста фич. Второе — нет: пережило ALS → CatBoost, S3 и расширение batch-полей. 

Именно поэтому гипотеза на живом трафике отработала — не идеальная ML-платформа, а правильная граница между данными и моделью.

Как вы запускали первый RecSys — с Feature Store или без? Что резали ради дедлайна?

Автор: LuckyChristi

Источник