Привет! Меня зовут Кристина, я 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]
}

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 для поиска отелей.

Рис. От оплаченного заказа к ответу API: target_order задаёт город и даты, Mongo отдаёт профиль и каталог, scoring ранжирует отели, наружу уходят готовые search_params
Рис. От оплаченного заказа к ответу API: target_order задаёт город и даты, Mongo отдаёт профиль и каталог, scoring ранжирует отели, наружу уходят готовые search_params

Фронту не отдаём сырые 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 merge_orders_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.

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

Один агрегат = один .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]
}

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

Результат

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

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

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

Так мы проверяли гипотезу в A/B: контроль — старая подборка, тест — Best Offer с моделью.

За месяц эксперимента конверсия из показа в заказ отеля выросла на единицы процентов, GMV - двузначный рост. Фичу раскатили на всех. Бэкенд за пять недель успел довести рекомендацию до экрана, где человек уже в контексте поездки.

Сервис отдаёт рекомендации на живом трафике с медианной задержкой ~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 или без? Что резали ради дедлайна?

Комментарии (0)