Привет, Хабр! Меня зовут Евстифеев Илья, я главный разработчик ML в команде MLOps FnR (Forecast & Replenishment) в MAGNIT TECH. Наша команда:

Ислам Янгуразов
Руководитель команды

Илья Евстифеев
Главный разработчик ML

Айрат Гибадуллин
Главный разработчик ML

Александр Щелоков
Ведущий разработчик ML

Михаил Проскурин
Ведущий разработчик ML

Денис Третьяков
Ведущий разработчик ML
Мы прогнозируем спрос и считаем пополнение для магазинов сети.
Сегодня я расскажу как наши ML пайплайны деплоит дежурный data scientist. Не ML-инженер, не DevOps — обычный DS. Инженеров данных в этом процессе нет не потому, что их нет, а потому что так задумано.
Я не буду пересказывать документацию Airflow, Spark или Delta Lake — вместо этого расскажу, как из их штатных возможностей у нас собрался CI/CD, в котором артефакт сборки — не бинарник, а код вместе с терабайтами рассчитанных таблиц.
Материал будет полезен инженерам данных, ML-инженерам, дата-сайентистам, которым надоело носить свои модели в прод через тикеты, — и всем, кто подозревает, что «MLOps-платформа» не обязана быть отдельным продуктом за отдельные деньги.
А чтобы было понятно, зачем всё это затевалось, — одна история. В феврале аналитик нашей команды решил проверить гипотезу: продажи блинной муки перед Масленицей растут не за неделю до праздника, как было заложено в прогнозе, а за две. Проверка заняла день — ноутбук, выборка по паре регионов, график, всё сходится. Дальше по классике должно было начаться самое долгое: тикет дата-инженерам, спека на таблицы, очередь, ревью, деплой — три-четыре недели, если у инженеров нет ничего срочного. Но мы знаем, что у DE всегда есть что-то срочное.
Вместо этого DS завёл в репозитории main.py на полсотни строк и config.yaml на двадцать, открыл MR — и через день после ревью расчёт поехал по расписанию на кластере в Kubernetes. DE в этой истории не появился ни разу. Вот как это устроено:

Как всё начиналось
Начиналось всё не с платформы и не с архитектурного комитета, а как у многих — с ноутбуков. Только не Jupyter, а Zeppelin: расчёты жили в ноутбуках и запускались ad-hoc — руками, когда понадобилось. Никакого расписания, никаких DAG'ов; «прод» — это ноутбук, который сегодня кто-то не забыл запустить.

Потом появился Airflow, а с ним первые DAG'и — казалось бы, вот она, взрослая жизнь. На деле мы просто перенесли старые привычки на новый инструмент: деплоили на горячую — поправил код, и следующим запуском он уже в проде, — а упавшие расчёты чинили руками прямо в бою. Вспоминать это больно: версий нет, откатываться некуда, стейджа не существует, а «работает» означало «пока никто не трогал».
Именно из этой боли, а не из красивой презентации, и выросли требования к нынешнему контуру. Хочется, чтобы расчёт нельзя было сломать правкой «на живую», — значит, релиз должен быть атомарным и версионированным. Хочется чинить не в бою — значит, нужны тестирование перед продом. Хочется, чтобы ошибка не превращалась в ночь героизма, — значит, нужен откат, который дешевле починки. Тогда мы ещё не знали, что самой интересной частью окажется не фремворк для джоб, а релизный цикл вокруг него, — но об этом дальше.

Два слова о фундаменте, на котором всё стоит. Все таблицы у нас — Delta Lake поверх S3, оркестрация — Airflow, расчёты — PySpark в Kubernetes. Формат таблиц мы не выбирали — Delta Lake стандарт платформы данных «Магнита», и о том, почему платформа в 2020-м поставила на него против Iceberg и Hudi и что из этого вышло за пять лет, коллеги уже написали: «Только Сигма выбирают Delta Lake». Пересказывать не будем; для нашей истории важно одно — формат таблиц оказался достаточно богатым, чтобы собрать на нём CI/CD для данных.
Что у нас получилось в цифрах
Прежде чем нырять в детали — масштаб системы, о которой пойдёт речь:
43 000 магазинов и 1 200 000 SKU, покрываемых прогнозом
один релиз - 200 Тб данных
offline-пайплайн — 45 задач с зависимостями: DQ-проверки источников, сезонность, эластичность, фактические цены, итоговый датасет и т.д.;
пять репозиториев участвуют в релизе: оффлайн-модели, инференс, генерация плана прогноза, интеграция, DQ;
от одного до шести релизов в месяц — катит дежурный DS;
Теперь по порядку: сначала покажем, как выглядит одна джоба глазами DS, потом пройдём по релизному циклу, расскажем историю одного отката и честно перечислим, что осталось ручным и на какие грабли мы наступили.
Как выглядит типичное приложение на нашем фреймворке PySpark Toolkit?
Разберём на живом примере — упрощённой версии нашей продовой джобы sales_uplift.
Со стороны DS джоба — это два файла. Один из них, main.py:
from dataclasses import dataclass import datetime from pyspark.sql import SparkSession, functions as F from pyspark_toolkit.configuration import JobConfiguration from pyspark_toolkit.configuration.schemas import Inputs from pyspark_toolkit.data_provider.delta import DeltaDataStorageClient, Merge from pyspark_toolkit.runner import spark_entrypoint from pyspark_toolkit.logger import pipeline_logger as logger @dataclass(frozen=True) class UpliftInputs(Inputs): end_date: datetime.date window_days: int @spark_entrypoint(config_input_type=UpliftInputs) def main(spark: SparkSession, data_storage_client: DeltaDataStorageClient, config: JobConfiguration): sales = data_storage_client.read_external_table(alias='sales_daily') holidays = data_storage_client.read_table( table_name=config.inputs.input_tables['holidays'], ) logger.info(f'считаем аплифт до {config.inputs.end_date}') uplift = ( sales .where(F.col('date_id') <= F.lit(config.inputs.end_date)) .join(holidays, on='date_id') # ... собственно расчёт аплифта ) data_storage_client.write_table( table_name=config.outputs.output_tables['sales_uplift'], df=uplift, write_options=Merge(merge_keys=['store_id', 'item_id', 'holiday_id']), ) if __name__ == '__main__': main()
Что здесь важно — не то, что есть, а то, чего нет:
нет создания SparkSession и разбора спарковых конфигов;
нет ни одного пути к S3, ни одного секрета, ни одной строки про Delta-обвязку;
нет argparse и ручного чтения YAML.
Функция получает готовую сессию, клиента хранилища и типизированный конфиг. UpliftInputs — обычный dataclass: если в YAML опечатка в имени поля или строка вместо даты, джоба упадёт на старте с внятной ошибкой валидации, а не через сорок минут внутри джойна.
При этом Spark мы сознательно не прятали. DS пишет обычный PySpark — джойны, оконки, репартиционирование остаются на его совести. Это принципиальный выбор, к нему вернёмся в разделе про грабли.
А теперь напишем config.yaml:
spark_conf: spark.driver.maxResultSize: 1g inputs: s3_bucket: ${S3_MAIN_BUCKET_NAME} end_date: 2026-07-01 window_days: 14 input_tables: holidays: holidays_description external_table_sources: sales_daily: table_name: sales schema: forecast type: hms unmanaged_external_table_sources: sales_daily: table_name: sales_daily_agg schema: dm type: hms outputs: output_tables: sales_uplift: sales_uplift secrets: aws_access_key: ${AWS_ACCESS_KEY} aws_secret_key: ${AWS_SECRET_KEY}
Параметр |
Описание |
|
Кастомный параметр |
|
Параметры Spark приложения |
|
Словарь с описанием таблиц, которые будут использоваться как входные данные. |
|
Словарь с описанием внешних таблиц, которые будут использоваться как входные данные. Под внешними таблицами здесь понимаются таблицы, которые не были напрямую порождены при помощи DataStorageClient из data_provider в рамках spark-приложения из текущего проекта. Важно отметить, что схемы таких таблиц могут меняться при помощи добавления префикса(dev_schema, stg_schema, dev_schema_branch) или постфикса независимо от пользователя при применении spark-toolkit-operator. |
|
Словарь с описанием внешних таблиц, которые имеют постоянную схему независимо от манипуляций с окружениями. С точки зрения самого тулкита физической разницы между unmanaged и managed нет - они все являются external и должны иметь уникальные ключи. |
|
Словарь с описанием выходных таблиц. Ключ - alias таблицы, значение - имя таблицы в файловой системе DataStorageClient. |
|
Здесь указываются чувствительные данные (у нас есть интеграция с Vault) |
Spark приложение становится таской в Airflow
Для запуска приложений мы написали свой оператор, наследовавшись от SparkKubernetesOperator.
Укажем название джобы из репозитория:

Чтобы создать таску нужно указать имя джобы, тег репозитория, параметры — всё. Оператор сам найдёт джобу по имени, поднимет под в K8s с нужным образом, подставит конфиг и секреты:
mk_holidays_sales_uplift = SparkToolKitKubernetesOperator( job_name="sales_uplift", repo_name="offline-pipeline", repo_branch=CURRENT_OFFLINE_TAG, # из Airflow Variable job_args={"end_date": "{{ params.end_date }}"}, )

Но таск в DAG'е — это ещё не прод. Прод — это когда новая версия модели доезжает до расписания, не сломав вчерашний прогноз, а если сломала — откатывается за минуты. И вот тут начинается самое интересное.
Единица релиза — тег, и данные версионируются вместе с кодом
Каждый релиз offline-части — это git-тег.

Все таблицы, которые считает пайплайн, живут в S3 по пути, включающему версию: <репозиторий>/<версия>. Есть тонкость, которую мы подсмотрели у семантического версионирования (и LakeFS). Оператор нормализует тег до минорной версии: v3.4.1 превращается в 3.4.
Тип релиза |
Сценарий |
Мажорный |
Расчет с нуля |
Минорный |
Релиз, основой которого является другой релиз |
Хотфикс |
Расчет "на горячую" |

У нас живет домовой

У веток данных есть и обратная сторона git-аналогии: как ветки в репозитории копятся, так и данные релизов никто добровольно не удаляет. Терабайты, в отличие от веток, стоят денег, поэтому мы написали keephouse — сервис-уборщик, который удаляет нерелизные таблицы и данные фича веток после 90 дней без записи — активность он определяет по истории Delta-таблицы. Релизные данные он не трогает: возможность отката — это то, ради чего всё затевалось.
Окружения — это линковка схем, а не копии данных
Физически таблица одна — лежит в S3 в каталоге своей версии. А вот имя, под которым её видят потребители, назначает linker — отдельная джоба на том же тулките, из ops-репозитория.

Она сканирует Delta-таблицы в каталоге <репозиторий>/<версия> и регистрирует их в Hive Metastore в схеме с префиксом окружения: dev, stg, uat, prod или без префикса.
Примеры таких схем:
Среда |
Таблица в HMS |
Путь |
dev |
dev_dm.table_name |
dev/offline-pipeline/feature_branch/table_name |
prod |
dm.table_name |
prod/v1.0/offline-pipeline/table_name |
stg |
stg_dm.table_name |
prod/v2.0/offline-pipeline/table_name |
uat |
uat_dm.table_name |
uat/v1.0/offline-pipeline/table_name |
При желании к имени таблицы добавляется постфикс с версией — так рядом могут сосуществовать таблицы нескольких релизов.
Для DS это выглядит как магия: в его конфиге написано dm, а в какую физическую схему это превратится — решает выбранное окружение перед запуском DAG.
Фича-ветка - это эксперимент
Тот же механизм «каталог в S3 + linker» даёт то, ради чего в классической разработке поднимают preview-окружения на каждый PR, только для данных.
Хочешь проверить гипотезу в ветке? DEV DAG сделает shallow clone данных текущего прод-тега в каталог твоей ветки. Это штатный SHALLOW CLONE из Delta: у новой таблицы свои метаданные, а файлы данных — общие с источником, поэтому клон терабайтной таблицы занимает секунды. Дальше DAG прогонит нужные задачи пайплайна уже в этом каталоге и залинкует результат в dev- схемы. Прод при этом не видит твою ветку в принципе: у неё свой каталог и свои схемы.

Корпоративные витрины (dm, dds_*) в dev-бакет тоже попадают не через поход в продовые таблицы напрямую. Отдельный ежедневный OPS DAG делает их shallow clone из prod- и uat-бакетов в dev-бакет и линкует в dev-схемы. Так эксперименты гоняются на свежих данных DLH, но на изолированной копии — прод об этом не знает.

DEV DAG — одно из редких мест, где участвует человек: здесь запускается эксперимент, результатом которого является отчет по метрикам.

Еще одна фича платформы: multibranch-DAG. В форме выбираешь пары веток inference#offline — хоть несколько — и Airflow через dynamic task mapping запускает DEV DAG по каждой паре параллельно. Сравнение двух версий модели на одинаковых данных превращается из недельной сборки экспериментов в нажатие несколько кнопок в конфигурации дага.

Для grid search (перебора большого количества параметров в рамках одной фичаветки) есть ещё изоляция по s3 ключу c experiment_id — данные эксперимента уезжают в отдельный подкаталог experiments/<experiment_id>/....

Деплой — это тоже DAG
Самое непривычное для человека из классической разработки: у нас нет Jenkins/GitLab-пайплайна для релиза. Деплой — это Airflow DAG с формой параметров, и это осознанный выбор: релиз ML-системы — это на 90% оркестрация расчётов, а оркестратор расчётов у нас уже есть.

Выглядит это так. Релизы катит дежурный DS — не выделенный релиз-инженер. Дежурный открывает DAG и заполняет форму:
выпадающее меню с git-тегами по каждому из репозиториев — offline-модели, inference, integration, data quality;
поле «комментарий к деплою: что катим и зачем»;
флажки «пропустить offline-часть» (если катим только inference);
«начать с задачи N» (если прошлый прогон упал на середине и пересчитывать сначала не нужно).
Жмём Trigger.
Дальше — автоматика:

В корпоративный мессенджер команды уходит сообщение с полной сводкой тегов и комментарием: деплой начался, @channel.
При необходимости — deep clone данных с предыдущего тега. Он копирует данные предыдущей версии, чтобы продолжить писать инкременты уже в новом каталоге. Поскольку в OSS Delta Lake нет
DEEP CLONE- копирует все утилита s5cmd. Вот цифра, которой мы сами удивились: 50 ТБ в час на поде с 4 ядрами и 16 ГБ памяти, ресурсы обычного недорогого ноутбука. Узкое место здесь не CPU и не память, а способность S3 держать тысячи параллельных запросов, и s5cmd выжимает из этого максимум. Получаются ветки данных — почти как в git, только весят терабайты.Прогоняется offline-пайплайн.
Linker регистрирует стейджинговые таблицы, и на них запускается тестовый inference — полный прогноз на реальных данных, но в изолированных схемах.
Если всё живо — таблицы перелинковываются в боевые схемы, тег релиза записывается в лог и в Airflow Variable.
-
В корпоративный мессенджер уходит «накат завершён» — с той же сводкой.

Дальше - прод
Даги, которые запускаются по расписанию называются schedule (регулярные).
Дообучение, инференс, прогноз по расписанию про деплой не знают ничего. Они читают текущий тег из Variable и следующим запуском по расписанию просто подхватывают новый релиз.


Аудит на стеройдах
Одна из киллерфич для аудита - при каждой записи в таблицу в метаданные коммита Delta Lake сохраняются:

Структура репозитория DAG'ов
Структура папок с DAG'ами описывает релизный процесс, а не бизнес-домен. Бизнес домен описан уже в самих Spark приложениях, что позволяет абстрагироваться от бизнес-логики. Чего не скажешь про предыдущие версии структуры.

История одного отката

Однажды мы выкатили релиз, расширявший таблицу прогноза новыми полями. Прогнозом пользуется соседний доменный стрим — команда пополнения: F&R в Магните живёт двумя доменными стримами, «Прогнозирование» и «Пополнение», у каждого свой релизный цикл и свои приоритеты (как устроены команды — отдельная статья коллег). Обновление интеграции под новую схему таблицы по времени разъехалось с нашим релизом — классическая история про две команды и один контракт данных: по отдельности обе сделали всё правильно, а подвёл стык.
Дальше всё случилось ровно так, как было спроектировано. Дежурный вернул Variable на предыдущий тег — и регулярные DAG’и следующим запуском начали запускать предыдущий коммит, где таблица прогноза лежала в старой схеме. Пополнение продолжило получать привычные данные. Мы синхронизировались с коллегами, дождались их релиза и накатили свой пайплайн заново.
Именно после таких дней перестаёшь считать версионирование данных перестраховкой: откат кода без отката данных здесь не спас бы — новые поля уже лежали в таблицах.
Сравним с классикой
Мы собрали процесс, который решает это, — и когда закончили, заметили, что переизобрели классический CI/CD, только для данных и моделей:
В классическом CI/CD |
У нас |
Артефакт сборки |
git-тег + все таблицы, рассчитанные этим тегом |
Registry |
S3-ключ вида |
Staging-окружение |
|
Preview-окружение на PR |
shallow clone production + |
Интеграционные тесты |
Прогон offline-пайплайна на стейдже + тестовый inference |
Оценка качества |
DQ-джобы и оценка метрик |
Выкат в прод |
Перелинковка схем в HMS + переключение Airflow Variable с текущим тегом |
Rollback |
Возврат Airflow Variable на предыдущий тег — его данные никуда не делись |
Release notes |
Комментарий к деплою + сводка тегов, автоматически в корпоративный мессенджер |
Кнопка Deploy |
Trigger деплойного DAG'а с формой параметров |
Если бы мы деплоили обычный сервис, всё было бы понятно: собрали образ, прогнали тесты, выкатили на стейдж, посмотрели, выкатили в прод.
Проблема ML-пайплайна в том, что его «артефакт» — это не бинарник. Это код плюс терабайты рассчитанных таблиц и обученные модели. Выкатить новую версию кода без её данных — бессмысленно; пересчитать все данные с нуля на каждый релиз — расточительно; а сломанный релиз должен откатываться вместе с данными, а не только с кодом.
Что осталось ручным, недоделанным и на какие грабли мы наступили
Честности ради — список того, что не решено или решено осознанным компромиссом.
Выбор тегов и нажатие кнопки — ручные. И это не тех. долг, а by design. Потому что e2e тестом оффлайн пайлпайна является инференс пайплайн и их нельзя катить отдельно.
Ночные падения разбираются утром. Классических дежурств с ночными подъёмами нет: первую линию держит техподдержка, а упавшие ночью расчёты разбираются утром вместе с DS — прогноз считается заранее, и у пайплайна есть запас по времени на такой разбор.
Порог входа и человеческий фактор. Деплойный DAG предполагает, что ты понимаешь бизнес-домен, линковку схем, — форма с десятком параметров не выглядит простой. Первый релиз новичок катит в паре с тем, кто уже делал это много раз. Для частичного решения этой проблемы мы планируем сделать dag builder для нашей релизной схемы (yaml манифест, который будет генерировать DAG'и).
Фреймворк снимает рутину, но не думает за тебя. Мы сознательно не прятали Spark — и это значит, что кривой джойн остается на совести автора. Декоратор не спасёт от плохого репартиционирования; спасёт только ревью/опыт/LLM.
Что мы поняли за 2 года
Главный вывод простой. CI/CD для данных нельзя купить как продукт, и это не серебряная пуля — это свойство, которое появляется, когда сделаны три скучные вещи: данные версионируются вместе с кодом, окружения — это метаданные, а не копии, и у релиза есть тесты и откаты. Delta Lake, Airflow и S3 сами по себе этого не дают: без инженерной обвязки вокруг это просто сырая технология с сырой документацией. Почти всё, на чём держится наш релизный цикл, это штатные фичи open-source технологий, использованные под нужным углом.
Второй вывод — про людей. Кнопку жмёт DS не ради экономии на инженерах, а потому что тот, кто понимает данные, способен довезти её до прода быстро, с лучшими показателями по метрикам моделей. Всё — toolkit, linker, пайплайны — существуют ради этого.
И третий. Проект без возможности отката релизов — это пороховая бочка. Позаботьтесь об этом как можно раньше.