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

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

Фронту не отдаём сырые 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 или без? Что резали ради дедлайна?