MLOps Engineer: Подготовка к техническому интервью

Интенсивный курс по проектированию и внедрению ML-сервисов в промышленную эксплуатацию. Фокус на инженерных аспектах: CI/CD, Kubernetes, оптимизации инференса и архитектуре современных LLM-систем.

Жизненный цикл ML-модели в продакшене

Жизненный цикл ML-модели в продакшене

По статистике, до 80% успешных в лаборатории ML-моделей никогда не добираются до реальных пользователей. Представьте: Data Scientist создал отличную модель кредитного скоринга для малого бизнеса (КМБ). В Jupyter Notebook на исторических данных она показывает точность 95%. Но когда ее пытаются внедрить в банковскую систему, выясняется, что она требует слишком много памяти, не умеет обрабатывать пропущенные поля в реальном времени, а через месяц ее предсказания вообще перестают совпадать с реальностью. Почему так происходит? Потому что Jupyter Notebook — это не продукт.

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

Чем ML-система отличается от обычного ПО?

В классической разработке программного обеспечения (Software Engineering) поведение системы определяется исключительно кодом. Если вы написали функцию сортировки, она будет работать одинаково и сегодня, и через год.

В машинном обучении поведение системы определяется кодом и данными. Формула успеха меняется: System=Code+DataSystem = Code + Data.

Если данные в реальном мире меняются (а они меняются всегда), система начинает ошибаться, даже если в коде не изменилось ни единого байта. Из-за этого классический процесс CI/CD (Continuous Integration / Continuous Deployment) в MLOps расширяется новым понятием — CT (Continuous Training), то есть непрерывным переобучением.

Бесконечная петля MLOps

Жизненный цикл ML-модели — это не прямая линия от идеи до релиза, а бесконечный цикл. На собеседовании важно показать, что вы видите всю картину целиком.

Этот цикл можно разделить на четыре крупных блока:

  1. Управление данными (Data Engineering) Модель нужно кормить данными. На этом этапе сырые данные очищаются и превращаются в признаки (features). В enterprise-решениях для этого используют Feature Store — специализированные хранилища, которые гарантируют, что модель при обучении и при работе в продакшене использует одни и те же алгоритмы расчета признаков.
  2. Разработка и эксперименты (ML Development) Data Scientist обучает десятки вариантов моделей, меняя гиперпараметры. Чтобы не запутаться, MLOps-инженер предоставляет инструменты для трекинга экспериментов (например, MLflow). Лучшая модель сохраняется в Model Registry — версионированное хранилище, где у каждой модели есть статус (Staging, Production, Archived).
  3. Развертывание (Deployment & Serving) Модель извлекается из хранилища и оборачивается в сервис. Это может быть REST API для обработки запросов в реальном времени (Inference) или пакетная обработка (Batch), когда модель скорит миллионы клиентов ночью по расписанию. Здесь в игру вступают контейнеры и оркестрация.
  4. Мониторинг (Monitoring) Сервис запущен. Теперь мы должны следить не только за «железными» метриками (CPU, RAM, время ответа), но и за ML-метриками (распределение входных данных, точность предсказаний). Как только метрики падают — цикл запускается заново.

MLOps-инженер не обучает нейросети. Он строит автоматизированный завод, на котором эти нейросети собираются, тестируются, доставляются до клиента и отправляются на ремонт, когда ломаются.

Почему модели ломаются: Concept Drift и Data Drift

Главная причина, по которой цикл MLOps должен быть замкнутым — это деградация модели с течением времени. В контексте банковских моделей для корпоративного сегмента (КСБ) и малого бизнеса (КМБ), это особенно критично. Экономика меняется, и вчерашние паттерны перестают работать.

Интервьюеры любят спрашивать про разницу между двумя типами деградации:

  • Data Drift (Сдвиг данных). Изменилось распределение входных признаков, но сама суть целевой переменной осталась прежней. Пример: Банк запустил агрессивную рекламу кредитов для IT-компаний. Раньше заявки подавали в основном торговые точки, а теперь — IT-сектор. Модель не видела столько IT-компаний при обучении и начинает ошибаться, потому что входные данные стали другими.
  • Concept Drift (Сдвиг концепции). Изменилась сама связь между признаками и целевой переменной. То, что раньше было нормой, теперь стало аномалией. Пример: Центробанк резко поднял ключевую ставку. Теперь даже компании с отличной историей (которые модель считает надежными) начинают банкротиться из-за дорогих кредитов. Правила игры изменились.

Чтобы бороться с деградацией, MLOps-инженер настраивает триггеры. Когда система мониторинга замечает, что распределение данных сильно отклонилось от эталонного (на котором модель обучалась), автоматически запускается пайплайн переобучения (CT).

Резюме для интервью

Когда на собеседовании вас просят спроектировать ML-систему (ML System Design), никогда не начинайте с выбора архитектуры нейросети. Начинайте с жизненного цикла:

  1. Откуда мы берем данные и как их храним?
  2. Как мы версионируем эксперименты?
  3. Как мы упаковываем модель для инференса?
  4. Как мы поймем, что модель начала деградировать в продакшене, и как мы ее обновим?

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

Контейнеризация моделей с помощью Docker

Контейнеризация моделей с помощью Docker

Вы забрали обученную модель кредитного скоринга из Model Registry, где она получила статус «Production». Data Scientist передает вам файл с весами и Python-скрипт. Вы запускаете код на боевом сервере, и он падает с ошибкой: несовпадение версий scikit-learn или отсутствие нужной версии CUDA.

Это классическая проблема «на моем компьютере всё работало». В машинном обучении она стоит особенно остро, потому что ML-модели тянут за собой тяжелый и хрупкий шлейф зависимостей: от конкретных версий Python-библиотек до системных C-компиляторов и драйверов видеокарт.

Чтобы разорвать эту зависимость от окружения, мы используем контейнеризацию.

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

Анатомия ML-контейнера

В мире Docker всё начинается с инструкций, записанных в Dockerfile. На собеседованиях по MLOps вас обязательно попросят написать или оптимизировать такой файл.

Давайте посмотрим на правильный базовый Dockerfile для нашего сервиса кредитного скоринга:

# 1. Выбор базового образа
FROM python:3.10-slim

# 2. Настройка рабочей директории
WORKDIR /app

# 3. Установка системных зависимостей (если нужны)
RUN apt-get update && apt-get install -y --no-install-recommends \
    build-essential \
    && rm -rf /var/lib/apt/lists/*

# 4. Кэширование Python-библиотек
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt

# 5. Копирование кода и модели
COPY src/ ./src/
COPY models/scoring_v1.pkl ./models/

# 6. Безопасность: запуск от имени непривилегированного пользователя
RUN useradd -m mluser
USER mluser

# 7. Точка входа
CMD ["python", "src/inference.py"]

Этот файл читается сверху вниз. Каждая команда создает новый слой (layer) в итоговом образе. И здесь кроется главный секрет оптимизации, о котором спрашивают на технических интервью.

Оптимизация сборки: кэширование слоев

Обратите внимание на шаги 4 и 5. Почему мы сначала копируем только requirements.txt, устанавливаем зависимости, и лишь потом копируем папку src/ с кодом?

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

Код в папке src/ (например, логика обработки признаков) меняется разработчиками постоянно. А вот список библиотек в requirements.txt — редко. Если бы мы написали COPY . . (скопировать всё сразу) перед RUN pip install, то любое, даже самое мелкое изменение в коде инвалидировало бы кэш, заставляя Docker заново скачивать и устанавливать тяжелые ML-библиотеки (которые могут весить гигабайты). Разделение этих шагов сокращает время сборки CI/CD пайплайна с 10 минут до 10 секунд.

Борьба с размером образа

ML-образы имеют тенденцию бесконтрольно разрастаться. Образ с PyTorch или TensorFlow легко может превысить 3 ГБ3 \text{ ГБ}. Большие образы медленно скачиваются на серверы, что критично при автомасштабировании (когда нужно срочно поднять новые копии сервиса при наплыве пользователей).

Как инженер, вы должны знать правила диеты для контейнеров:

  1. Правильный базовый образ. Никогда не используйте полные образы вроде python:3.10 (они включают кучу лишних утилит). Используйте python:3.10-slim. Важное примечание для интервью: часто для уменьшения размера советуют образы на базе alpine. Для веб-разработки это отлично, но для ML — табу. Alpine использует библиотеку musl вместо стандартной glibc. Большинство ML-библиотек (numpy, pandas) скомпилированы под glibc. В Alpine они будут компилироваться из исходников при установке, что займет часы и приведет к багам.
  2. Флаг --no-cache-dir. При установке через pip всегда добавляйте этот флаг, чтобы не сохранять установочные архивы .whl внутри образа.
  3. Очистка системного кэша. Если вы устанавливали системные пакеты через apt-get, завершайте команду очисткой списков: rm -rf /var/lib/apt/lists/*.

Архитектурный выбор: где хранить веса модели?

В нашем примере мы скопировали файл scoring_v1.pkl прямо внутрь Docker-образа (шаг 5). Это архитектурное решение, у которого есть альтернатива. На собеседовании вас спросят: «Как лучше доставлять модель в контейнер?»

Существует два основных подхода:

Характеристика Вшивание модели (Bake-in) Загрузка при старте (Runtime Download)
Суть Файл модели копируется в образ на этапе сборки (docker build). Образ содержит только код. Модель скачивается из S3 / Model Registry при запуске контейнера.
Плюсы Абсолютная воспроизводимость. Образ самодостаточен. Быстрый старт контейнера. Образ весит мало. Можно обновить модель без пересборки кода (просто поменяв конфиг).
Минусы Образ становится огромным. При каждом переобучении (Continuous Training) нужно собирать новый Docker-образ. Контейнер стартует медленнее (ждет скачивания). Зависимость от сети и хранилища (S3 упал — сервис не поднялся).

Что выбрать? Если ваша модель весит мало (например, логистическая регрессия или случайный лес для КМБ на 50 МБ) — смело вшивайте её в образ. Это надежнее. Если вы работаете с Large Language Models (LLM), веса которых достигают десятков гигабайт, вшивать их в образ безумие. В таких случаях используют загрузку при старте или монтируют внешнее хранилище прямо в контейнер.

Итог

Мы упаковали наш сервис кредитного скоринга в легковесный, безопасный и воспроизводимый Docker-контейнер. Теперь он гарантированно будет работать одинаково и на ноутбуке разработчика, и на тестовом стенде.

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

Оркестрация контейнеров в Kubernetes для ML

Оркестрация контейнеров в Kubernetes для ML

Упакованная в Docker-образ модель кредитного скоринга успешно работает на сервере и выдает предсказания. Но что произойдет, если отдел маркетинга запустит масштабную кампанию, и поток заявок от малого бизнеса вырастет в 100 раз? Сервер исчерпает ресурсы, процесс упадет с ошибкой нехватки памяти, и бизнес начнет терять деньги с каждой минутой простоя. Docker отлично решает проблему изоляции зависимостей, но он не умеет управлять кластером серверов, балансировать нагрузку и автоматически перезапускать упавшие контейнеры. Для этого нужен дирижер — оркестратор.

В enterprise-разработке стандартом де-факто для оркестрации стал Kubernetes (K8s). На техническом интервью MLOps-инженера ожидают не просто знания базовых команд K8s, а понимания того, как адаптировать его механизмы под специфику тяжелых ML-нагрузок.

От контейнера к кластеру: базовые примитивы

Kubernetes не работает с Docker-контейнерами напрямую. Он вводит собственный уровень абстракций, чтобы управлять инфраструктурой декларативно: мы описываем желаемое состояние системы, а K8s сам решает, как его достичь.

Pod (Под) — минимальная единица развертывания в K8s. Это «обертка» вокруг одного или нескольких тесно связанных контейнеров, которые делят общую память и сеть.

В контексте ML в одном Поде обычно крутится контейнер с моделью (инференс). Иногда к нему добавляют вспомогательный контейнер (sidecar), например, для сбора специфических метрик распределения данных (Data Drift) и отправки их в систему мониторинга.

Поды смертны. Если сервер (Node), на котором запущен Под, выходит из строя, Под исчезает навсегда. Чтобы обеспечить отказоустойчивость, K8s использует контроллеры:

  1. Deployment (Развертывание) — следит за тем, чтобы в кластере всегда работало заданное количество копий (реплик) Пода. Если Под падает, Deployment немедленно создает новый на здоровом узле.
  2. Service (Сервис) — обеспечивает единую точку входа. Каждый раз, когда Deployment пересоздает Под, тот получает новый внутренний IP-адрес. Service дает стабильный IP и балансирует входящие запросы к модели между всеми живыми репликами.

Автомасштабирование инференса (HPA)

Держать 100 Подов с моделью постоянно запущенными — дорого. Инфраструктура должна адаптироваться под нагрузку. За это отвечает Horizontal Pod Autoscaler (HPA).

HPA постоянно опрашивает метрики потребления ресурсов (например, CPU) и динамически меняет количество реплик в Deployment. Логика масштабирования описывается следующей формулой:

DesiredReplicas=CurrentReplicas×CurrentMetricValueDesiredMetricValueDesiredReplicas = \left\lceil CurrentReplicas \times \frac{CurrentMetricValue}{DesiredMetricValue} \right\rceil

Где:

  • DesiredReplicas\mathbf{DesiredReplicas} — целевое количество реплик (округленное вверх \lceil \dots \rceil).
  • CurrentReplicas\mathbf{CurrentReplicas} — текущее количество работающих Подов.
  • CurrentMetricValue\mathbf{CurrentMetricValue} — текущее потребление ресурса (например, 80% CPU).
  • DesiredMetricValue\mathbf{DesiredMetricValue} — желаемое потребление, заданное инженером (например, 50% CPU).

Практический пример: У вас работают 2 реплики модели скоринга. Из-за наплыва заявок средняя загрузка CPU достигла 80%. В настройках HPA указан таргет в 50%. Считаем: 2×(80/50)=3.22 \times (80 / 50) = 3.2. K8s округляет это значение вверх до 4. HPA дает команду Deployment увеличить количество реплик до четырех, тем самым размазывая нагрузку и снижая CPU на каждом Поде до целевых значений.

Управление ресурсами: Requests и Limits

ML-модели крайне требовательны к оперативной памяти (RAM) при загрузке весов и к CPU/GPU при инференсе. Если не ограничить аппетиты Пода, одна тяжелая модель может захватить все ресурсы узла, "убив" соседние сервисы.

В манифестах K8s ресурсы описываются двумя параметрами:

Параметр Смысл для планировщика K8s Что будет при превышении?
Requests (Запросы) Гарантированный минимум. K8s ищет сервер, где есть столько свободного места, чтобы разместить Под. Если запрошено больше, чем есть в кластере — Под зависнет в статусе Pending.
Limits (Лимиты) Жесткий потолок. Максимум ресурсов, который разрешено потребить Поду. CPU: процесс будет искусственно замедлен (троттлинг).<br>RAM: Под будет убит системой с ошибкой OOMKilled (Out Of Memory).

Частая ошибка на проде: Размер весов модели скоринга на диске — 2 ГБ. Инженер ставит limits.memory: 2.5Gi. При старте Python-процесс десериализует модель в оперативную память, потребление скачет до 3.5 ГБ. K8s моментально убивает Под (OOMKilled). Модель уходит в бесконечный цикл перезапусков (CrashLoopBackOff). Правило: Лимиты памяти должны учитывать пиковое потребление при загрузке модели и обработке батча данных, а не только размер файла весов.

Специфика ML: маршрутизация на GPU-узлы

Кластер обычно состоит из множества серверов (Nodes), но видеокарты (GPU) установлены лишь на некоторых из них. GPU — дорогой ресурс. Нам нужно гарантировать две вещи:

  1. Тяжелые модели (например, нейросети для обработки документов клиентов) должны запускаться только на узлах с GPU.
  2. Обычные легковесные микросервисы (например, логирование) не должны занимать место на GPU-узлах.

В K8s это решается механизмом Taints and Tolerations (Отторжения и Допущения) в связке с nodeSelector.

На GPU-сервер вешается «отторжение» (Taint): «Сюда нельзя селиться никому, кроме тех, у кого есть специальный допуск». Обычные Поды будут обходить этот сервер стороной. В манифесте Пода с тяжелой ML-моделью мы прописываем «допущение» (Toleration) к этому Taint, а также nodeSelector, который прямо указывает K8s: «Ищи сервер с меткой hardware=gpu».

Бесшовное обновление моделей (Rolling Update)

Жизненный цикл ML-системы подразумевает регулярное переобучение (Continuous Training). Когда новая версия модели упакована в Docker-образ, ее нужно выкатить в продакшен без остановки сервиса.

Deployment в K8s по умолчанию использует стратегию Rolling Update. Вместо того чтобы убить все старые Поды и запустить новые (что вызвало бы даунтайм), K8s делает это постепенно:

  1. Создает один Под с новой версией модели.
  2. Ждет, пока Под сигнализирует, что модель загружена в память и готова принимать трафик (за это отвечают Readiness Probes — пробы готовности).
  3. Переключает часть трафика на новый Под.
  4. Убивает один старый Под.
  5. Повторяет процесс, пока все Поды не будут обновлены.

Если новая модель падает при старте (например, несовместимость версий библиотек), проба готовности провалится. K8s остановит обновление, оставив работать старые, проверенные реплики.

Обернув контейнер с моделью в манифесты Kubernetes, мы получили масштабируемую, самовосстанавливающуюся систему, способную пережить пиковые нагрузки и безопасно обновляться. Но сам контейнер пока остается "черным ящиком". Чтобы Service мог отправлять в него HTTP-запросы, а Readiness Probes — проверять статус, внутри контейнера должен работать быстрый и надежный веб-сервер.

Проектирование API для инференса на FastAPI

Проектирование API для инференса на FastAPI

Контейнер собран, лимиты памяти в Kubernetes настроены, Под запущен. Но внутри этого изолированного окружения лежит лишь сериализованный файл с весами модели. Чтобы другие микросервисы могли отправлять данные на скоринг, нам нужен мост между миром HTTP-запросов и математикой машинного обучения. Этим мостом выступает веб-сервер.

В современном Python-стеке стандартом де-факто для ML-инференса стал FastAPI. Он быстрый, автоматически генерирует документацию (Swagger) и из коробки решает две главные боли MLOps-инженера: валидацию запутанных входных данных и управление жизненным циклом приложения.

Управление состоянием: где должна жить модель?

Самая частая ошибка новичков при проектировании API — загрузка модели прямо внутри функции-обработчика запроса (эндпоинта). Если десериализация случайного леса занимает 500 миллисекунд, каждый клиент будет ждать эти полсекунды плюсом к самому инференсу.

Модель должна находиться в оперативной памяти постоянно. Но если мы просто объявим глобальную переменную на уровне модуля, загрузка начнется в момент импорта файла. Это усложняет тестирование и делает поведение приложения непредсказуемым при запуске через ASGI-серверы (например, Uvicorn).

Правильный паттерн в FastAPI — использование контекстного менеджера lifespan.

Lifespan — это механизм FastAPI, позволяющий выполнить тяжелую инициализацию (загрузку весов, подключение к БД) строго до того, как сервер начнет принимать HTTP-запросы, и корректно освободить ресурсы при остановке.

Именно здесь происходит стыковка с Kubernetes, о которой мы говорили в прошлой главе. K8s регулярно опрашивает наш сервис через Readiness Probe (пробу готовности). Если мы свяжем эндпоинт /health с успешным завершением блока инициализации lifespan, Kubernetes направит реальный пользовательский трафик на Под только тогда, когда модель полностью загружена в RAM.

Жесткие контракты данных с Pydantic

Модели машинного обучения хрупки. Классический код на Pandas или NumPy ожидает матрицу чисел. Если фронтенд пришлет возраст клиента в виде строки "двадцать" вместо числа 20, ML-модель упадет с нечитаемой ошибкой глубоко в недрах библиотеки C++.

API обязано выступать жестким фильтром. В FastAPI эта задача решается через Pydantic — библиотеку валидации данных на основе аннотаций типов Python.

Мы описываем ожидаемую структуру (схему) запроса. Если данные не соответствуют контракту, FastAPI даже не передаст их в функцию инференса, а мгновенно вернет клиенту HTTP-статус 422 Unprocessable Entity с подробным описанием ошибки.

from pydantic import BaseModel, Field

class ScoringRequest(BaseModel):
    age: int = Field(..., ge=18, le=100, description="Возраст заемщика")
    income: float = Field(..., gt=0, description="Ежемесячный доход в руб.")
    business_sector: str = Field(..., description="Отрасль КМБ")

class ScoringResponse(BaseModel):
    probability_of_default: float
    decision: str

Ловушка асинхронности: главный вопрос на собеседовании

FastAPI построен на асинхронной архитектуре (ASGI) и использует Event Loop (цикл событий). Это позволяет серверу держать тысячи одновременных соединений, переключаясь между ними, пока ожидается ответ от базы данных или сети (I/O-bound задачи).

Но инференс большинства классических ML-моделей (Scikit-learn, XGBoost, CatBoost) — это CPU-bound задача. Она непрерывно нагружает процессор математическими вычислениями и не поддерживает асинхронность (await).

Если вы вызовете синхронный метод model.predict() внутри асинхронного эндпоинта (async def), вы заблокируете Event Loop.

Представим, что время инференса одного запроса T=0.2T = 0.2 секунды. Максимальная пропускная способность (RPS) заблокированного воркера составит всего RPS=1/TRPS = 1 / T, то есть 5 запросов в секунду. Пока модель считает эти 0.2 секунды, сервер "зависнет": он не сможет принять новые соединения, и даже эндпоинт /health перестанет отвечать. Kubernetes решит, что Под "мертв", и перезапустит его.

Чтобы этого избежать, FastAPI применяет элегантное правило маршрутизации под капотом:

Сигнатура эндпоинта Что происходит под капотом Когда применять в ML
async def predict(...) Выполняется прямо в главном Event Loop. Блокирует его при долгих вычислениях. Если модель поддерживает await (редкость) или если API просто вызывает внешний микросервис/БД.
def predict(...) FastAPI автоматически отправляет эту функцию в внешний пул потоков (Threadpool). Event Loop остается свободным. Для 95% ML-моделей. Вычисления идут в отдельном потоке, сервер продолжает отвечать на /health.

Собираем production-ready API

Объединим управление состоянием, валидацию Pydantic и правильную работу с потоками в единый сервис для нашего кредитного скоринга:

from contextlib import asynccontextmanager
from fastapi import FastAPI, HTTPException
import joblib
import pandas as pd

# Глобальный словарь для хранения состояния приложения
ml_models = {}

@asynccontextmanager
async def lifespan(app: FastAPI):
    # Startup: загружаем модель в память до старта сервера
    ml_models["kmb_scorer"] = joblib.load("/models/kmb_scorer_v1.pkl")
    yield
    # Shutdown: очищаем ресурсы при остановке
    ml_models.clear()

app = FastAPI(lifespan=lifespan, title="KMB Scoring API")

@app.get("/health")
async def health_check():
    # Если lifespan не отработал, ключа не будет, вернется 503
    if "kmb_scorer" not in ml_models:
        raise HTTPException(status_code=503, detail="Model not loaded")
    return {"status": "ready"}

# Обратите внимание: используем def, а не async def для CPU-bound задачи!
@app.post("/predict", response_model=ScoringResponse)
def predict(request: ScoringRequest):
    model = ml_models["kmb_scorer"]

    # Преобразуем валидные данные Pydantic в формат для модели
    input_data = pd.DataFrame([request.model_dump()])

    # Синхронный вызов, который FastAPI выполнит в отдельном потоке
    prediction = model.predict_proba(input_data)[0][1]

    decision = "approve" if prediction < 0.15 else "decline"

    return ScoringResponse(
        probability_of_default=prediction,
        decision=decision
    )

Этот код полностью готов к интеграции с инфраструктурой. Kubernetes сможет проверять /health, Pydantic защитит от мусорных данных, а синхронный def убережет сервер от зависаний под нагрузкой.

Однако, даже если мы не блокируем сервер, сам по себе инференс в Python через .predict() может быть недостаточно быстрым для высоконагруженных систем или тяжелых нейросетей. В следующей главе мы разберем, как ускорить саму математику модели.

Оптимизация инференса и ускорение моделей

Оптимизация инференса и ускорение моделей

Если FastAPI-приложение из предыдущего этапа начинает захлебываться под нагрузкой, Kubernetes послушно поднимет новые реплики Подов через HPA. Но масштабирование инфраструктуры — это борьба с симптомами. Если сама ML-модель тратит на один прогноз кредитного рейтинга 500 миллисекунд, вы будете сжигать серверный бюджет на неэффективные вычисления.

Проблема кроется в том, что фреймворки вроде PyTorch или pandas созданы для удобства разработки и обучения, а не для миллисекундного продакшена. Чтобы ускорить инференс (процесс получения предсказаний), необходимо изменить формат модели, снизить математическую точность и перестроить логику подачи данных.

1. Компиляция графа: переход к ONNX

В момент вызова model.predict() в нативном Python-фреймворке происходит множество проверок типов и выделений памяти. Для инференса эта гибкость избыточна: архитектура модели уже зафиксирована.

Первый шаг к ускорению — экспорт модели в промежуточное представление, например, ONNX (Open Neural Network Exchange). ONNX превращает модель в статический вычислительный граф, который затем исполняется через специализированный движок — ONNX Runtime.

ONNX Runtime применяет графовые оптимизации. Например, если в нейронной сети за операцией умножения на матрицу весов сразу идет функция активации (ReLU), движок объединяет их в одну инструкцию (Operator Fusion). Это исключает промежуточную запись данных в память.

Перевод модели скоринга КМБ (XGBoost или PyTorch) в формат ONNX и запуск через ONNX Runtime обычно дает ускорение от 2x до 5x просто за счет устранения накладных расходов Python и слияния операций.

2. Квантование: снижение точности вычислений

По умолчанию веса нейронных сетей хранятся в формате 32-битных чисел с плавающей точкой (FP32). Однако для большинства ML-задач такая ювелирная точность не требуется.

Квантование — это процесс преобразования весов и активаций модели из формата высокой точности (FP32) в формат пониженной точности, чаще всего в 8-битные целые числа (INT8).

Потребление памяти MM для модели с NN параметрами напрямую зависит от выбранного типа данных:

M=N×BM = N \times B

где BB — количество байт на один параметр. Для FP32 B=4B = 4, а для INT8 B=1B = 1.

Если наша модель содержит 100 миллионов параметров, в FP32 она займет 400 МБ оперативной памяти, а в INT8 — всего 100 МБ. Но выигрыш заключается не только в экономии RAM.

Современные процессоры (CPU) тратят больше времени на перемещение данных из оперативной памяти в кэш процессора, чем на само умножение чисел. Сжав модель в 4 раза, мы в 4 раза ускоряем пропускную способность памяти (Memory Bandwidth). Кроме того, целочисленная арифметика на уровне тактов процессора выполняется быстрее, чем операции с плавающей точкой.

Процесс перевода числа из FP32 в INT8 опирается на вычисление масштабного коэффициента (Scale):

S=WmaxWmin2b1S = \frac{W_{max} - W_{min}}{2^b - 1}

Здесь WmaxW_{max} и WminW_{min} — максимальное и минимальное значения весов в слое, а bb — целевая битность (для INT8 знаменатель равен 255). Каждое число делится на SS и округляется до ближайшего целого.

Плата за квантование — незначительное падение метрик качества (например, ROC-AUC может снизиться на 0.005). В инженерной практике это допустимый компромисс ради кратного снижения задержки (latency).

3. Динамический батчинг (Dynamic Batching)

Матричные вычисления под капотом ML-моделей (особенно при использовании векторизации на CPU или ядер CUDA на GPU) устроены так, что обработать 10 запросов одновременно почти так же быстро, как обработать 1 запрос.

Однако в реальном API запросы от клиентов приходят не пачкой, а последовательно и асинхронно. Если отправлять каждый запрос в модель поштучно, вычислительные ресурсы будут простаивать.

Динамический батчинг решает эту проблему, создавая микро-очередь на стороне сервера инференса:

  1. Приходит первый запрос на скоринг. Сервер не отправляет его в модель сразу.
  2. Открывается окно ожидания (например, 10 миллисекунд).
  3. Если за эти 10 мс приходят еще 4 запроса, сервер объединяет их в один тензор (батч размером 5).
  4. Модель делает предсказание для всего батча за один проход.
  5. Сервер "распаковывает" результаты и отдает каждому клиенту его ответ.

Именно здесь раскрывается потенциал асинхронности (Event Loop), заложенный в FastAPI на предыдущем этапе. Пока формируется батч, сервер не блокируется и продолжает принимать новые соединения.

Сравнение подходов к оптимизации

Метод Что улучшает Главный плюс Главный минус
ONNX Runtime Вычислительный граф Работает "из коробки", не меняет точность предсказаний Не все экзотические слои PyTorch поддерживаются при экспорте
Квантование (INT8) Потребление памяти и скорость Снижает размер модели в 4 раза Требует валидации метрик качества (возможна деградация)
Динамический батчинг Пропускную способность (RPS) Максимально утилизирует CPU/GPU Увеличивает базовую задержку (latency) на время окна ожидания

Внедрение этих трех практик превращает тяжеловесный Python-скрипт в высокопроизводительный сервис. Однако оптимизированная модель будет работать быстро только в том случае, если данные для нее поставляются без задержек. Если инференс занимает 10 мс, а сбор признаков из базы данных — 2 секунды, вся оптимизация теряет смысл.

Хранение признаков и Feature Store в реальном времени

Хранение признаков и Feature Store в реальном времени

Мы перевели модель в формат ONNX, применили INT8-квантование и настроили динамический батчинг. Теперь сам инференс занимает ничтожные 5 миллисекунд. Но при нагрузочном тестировании FastAPI-сервиса мы видим, что полный ответ клиенту занимает 505 миллисекунд. Куда уходят еще полсекунды? Они тратятся на то, чтобы сходить в аналитическую базу данных, выполнить тяжелый SQL-запрос и собрать историю транзакций клиента для передачи в модель.

Время ответа системы описывается простой формулой: Ttotal=Tdata+TinferenceT_{total} = T_{data} + T_{inference} Где TtotalT_{total} — общее время ответа, TdataT_{data} — время подготовки данных, а TinferenceT_{inference} — время вычислений модели. Если TdataT_{data} на два порядка больше TinferenceT_{inference}, все усилия по ускорению вычислений теряют смысл. Нам нужно доставлять данные в модель так же быстро, как она их обрабатывает.

Проблема рассинхронизации: Training-Serving Skew

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

Когда дата-саентист обучает модель кредитного скоринга для малого бизнеса (КМБ), он работает с историческими данными. Он может запустить тяжелый PySpark-скрипт поверх кластера Hadoop, который будет несколько часов агрегировать транзакции миллионов клиентов за последние 5 лет, чтобы вычислить признак «среднемесячный оборот». Вектор оптимизации здесь — пропускная способность (Throughput) при работе с терабайтами данных.

Когда модель развернута в Kubernetes за FastAPI-сервисом, она обрабатывает запросы по одному (или небольшими батчами). Приходит запрос: {"client_id": "12345"}. Сервису нужно получить тот же самый «среднемесячный оборот» для этого конкретного клиента прямо сейчас. Вектор оптимизации здесь — минимальная задержка (Latency).

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

Training-Serving Skew — ситуация, при которой логика расчета признаков или сами данные на этапе обучения модели расходятся с тем, что модель получает на этапе реального инференса. Это приводит к непредсказуемой деградации качества прогнозов в продакшене.

Архитектура Feature Store

Чтобы разорвать зависимость между медленной аналитикой и быстрым инференсом, в MLOps архитектуру внедряют платформу управления признаками.

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

Под капотом Feature Store — это не одна база данных, а комбинация двух хранилищ с механизмом синхронизации между ними.

Характеристика Offline Store (Офлайн-хранилище) Online Store (Онлайн-хранилище)
Технологии Hadoop, S3, PostgreSQL, ClickHouse Redis, Memcached, DynamoDB
Объем данных Петабайты (вся история за годы) Гигабайты (только актуальные срезы)
Паттерн доступа Сканирование таблиц, джойны, батчи Чтение по ключу (Key-Value lookup)
Задержка (Latency) Минуты / Часы Миллисекунды (< 5 мс)
Потребители Обучение моделей (Continuous Training) FastAPI-сервисы (Inference)

Как это работает на практике

Представим систему кредитного скоринга. У нас есть признак avg_monthly_turnover_6m (среднемесячный оборот за 6 месяцев).

  1. Запись: Каждую ночь по расписанию запускается тяжелый процесс обработки данных. Он сканирует транзакции в Hadoop, высчитывает актуальный оборот для всех клиентов и сохраняет полную историю в Offline Store (для будущих переобучений).
  2. Синхронизация: Сразу после этого процесс берет только самые свежие значения этого признака для каждого клиента и записывает их в Online Store (Redis) в виде простой пары ключ-значение: {"client_12345:avg_monthly_turnover_6m": 1500000.00}.
  3. Чтение: Когда днем клиент подает заявку на кредит, FastAPI-сервис делает O(1) запрос в Redis по ключу клиента, мгновенно получает готовое число 1500000.00 и передает его в ONNX-модель.

Благодаря этому TdataT_{data} снижается с 500 мс до 1-2 мс.

Защита от заглядывания в будущее (Point-in-Time Correctness)

Feature Store решает еще одну критическую проблему подготовки данных — утечку данных (Data Leakage) при формировании обучающих выборок.

Представьте, что мы собираем датасет для обучения новой версии скоринговой модели. У нас есть клиент, который взял кредит 1 января 2023 года и допустил дефолт в июне. Мы хотим обучить модель предсказывать этот дефолт. Для этого нам нужно взять признаки клиента.

Если мы просто возьмем его текущее состояние из базы данных (например, «оборот за последние 6 месяцев», рассчитанный на сегодня), мы совершим фатальную ошибку. Мы покажем модели данные из будущего по отношению к моменту выдачи кредита. Модель выучит, что перед дефолтом оборот падает, но в реальности 1 января (в момент принятия решения) этого падения еще не было.

Feature Store решает это с помощью механизма Point-in-Time (PiT) Join.

Point-in-Time Join — алгоритм объединения таблиц, при котором для каждого целевого события (например, заявки на кредит) из истории извлекаются значения признаков, актуальные строго на момент времени до наступления этого события.

Когда мы запрашиваем датасет из Offline Store, мы передаем сущность и временную метку: (client_id="12345", timestamp="2023-01-01 10:00:00"). Feature Store автоматически отматывает историю изменений признака и возвращает то значение оборота, которое было известно системе ровно на утро 1 января, игнорируя все последующие транзакции. Это гарантирует, что Continuous Training пайплайн всегда обучается на честных исторических срезах.

Ограничения пакетного обновления

Схема с ночным расчетом признаков и выгрузкой их в Redis отлично работает для медленно меняющихся данных (оборот за полгода, возраст бизнеса, кредитный рейтинг).

Но что, если для защиты от фрода нам нужен признак count_transactions_10min (количество транзакций за последние 10 минут)? Ночной расчет здесь бессилен — данные устареют к утру. Нам нужно вычислять признаки на лету, прямо в момент совершения транзакций, и мгновенно обновлять Online Store. Для этого инфраструктура должна перейти от пакетной обработки к потоковой.

Потоковая обработка данных с Kafka и Spark Streaming

Потоковая обработка данных с Kafka и Spark Streaming

Представьте, что клиент подает заявку на кредит для бизнеса. Наш Online Feature Store (Redis) мгновенно отдает модели его средний оборот за полгода, и модель готова одобрить выдачу. Но есть нюанс: всего 30 секунд назад со счетов этого клиента начался подозрительный вывод средств. Если наш Feature Store обновляется ночным скриптом из базы данных, модель просто не увидит эту аномалию и пропустит фрод. Чтобы реагировать на действия пользователя «здесь и сейчас», признаки должны вычисляться в реальном времени.

Для решения этой задачи в MLOps используется связка систем потоковой передачи событий и движков вычислений в памяти.

Apache Kafka: Нервная система данных

Когда в системе происходит событие (транзакция, клик, смена геолокации), оно должно быть мгновенно доставлено потребителям. Эту роль выполняет Apache Kafka.

В вакансиях часто требуют опыт работы с очередями (RabbitMQ, Celery) и Kafka. Важно понимать фундаментальную разницу между ними на техническом интервью: Kafka — это не классическая очередь задач, это распределенный журнал событий (commit log).

Характеристика Классическая очередь (RabbitMQ, Celery) Брокер событий (Apache Kafka)
Жизненный цикл сообщения Удаляется после успешного прочтения (ACK). Сохраняется на диске заданное время (например, 7 дней).
Паттерн потребления Конкурирующие воркеры (сообщение достается одному). Множество независимых систем могут читать одни и те же данные.
Порядок сообщений Может нарушаться при повторных попытках (retries). Строго гарантирован внутри одной партиции.

В Kafka данные организованы в топики (Topics) — логические каналы для однородных событий, например, transactions_topic. Для масштабирования топик бьется на партиции (Partitions).

Партиция — это физический файл на диске сервера Kafka, куда события записываются строго друг за другом (append-only). Каждое событие получает свой порядковый номер — оффсет (Offset).

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

Spark Streaming: Превращение потока в признаки

Kafka отлично хранит и доставляет байты, но для ML-модели нужны агрегированные признаки (фичи). Нам нужно на лету считать метрики вроде «сумма транзакций за последние 5 минут». Здесь в игру вступает Apache Spark Streaming.

Spark Structured Streaming использует концепцию микро-батчинга (Micro-batching). Вместо того чтобы обрабатывать каждое событие поштучно (что создает огромный накладной расход на сетевые вызовы), Spark накапливает события из Kafka за короткий интервал (например, 500 миллисекунд), превращает их в небольшой DataFrame и выполняет над ним привычные SQL-подобные трансформации.

Оконные функции и проблема опоздавших данных

При расчете потоковых признаков мы оперируем временными окнами. Но в распределенных системах возникает проблема: время, когда событие произошло (Event Time), и время, когда оно дошло до сервера (Processing Time), часто не совпадают.

Клиент совершил покупку в 12:01, заходя в метро. Связь пропала, и телефон отправил событие на сервер только в 12:06, когда клиент вышел на улицу. Если Spark будет считать агрегации по Processing Time, транзакция попадет в окно «12:05–12:10», что исказит реальную картину поведения клиента.

Поэтому вычисления всегда строятся на базе Event Time. Но как долго Spark должен держать в оперативной памяти агрегацию для окна «12:00–12:05», ожидая опоздавшие данные? Бесконечно копить состояние нельзя — мы получим OOMKilled (исчерпание памяти).

Для решения этой проблемы используется механизм Watermark (Водяной знак). Это динамический порог времени, старше которого данные гарантированно отбрасываются.

Формула расчета водяного знака выглядит так:

Twatermark=max(Tevent)ΔT_{watermark} = \max(T_{event}) - \Delta

Где TwatermarkT_{watermark} — граница отсечения опоздавших данных, max(Tevent)\max(T_{event}) — максимальное время события, которое мы уже успели прочитать из потока, а Δ\Delta — допустимое время опоздания.

Практический пример: Мы задали допустимое время опоздания Δ=2\Delta = 2 минуты.

  1. Spark читает поток. Последнее пришедшее событие имеет метку времени 12:10 (max(Tevent)\max(T_{event})).
  2. Spark вычисляет водяной знак: 12:102 мин=12:0812:10 - 2 \text{ мин} = 12:08.
  3. Если после этого в систему прилетает транзакция из метро с меткой Event Time 12:05, Spark ее проигнорирует, так как 12:05<12:0812:05 < 12:08. Память для старых окон уже очищена.

Архитектура в сборе: от транзакции до инференса

Теперь мы можем замкнуть архитектуру, объединив знания из предыдущих этапов проектирования ML-системы:

  1. Генерация события: Банковский бэкенд (Producer) записывает JSON с деталями транзакции в топик Kafka client_transactions.
  2. Потоковая агрегация: Приложение на PySpark непрерывно читает этот топик. Оно группирует данные по client_id и скользящему окну в 10 минут (используя Event Time и Watermark).
  3. Обновление Online Store: Вычисленный признак (например, tx_sum_10m) Spark Streaming немедленно записывает (Upsert) в Redis.
  4. Инференс: Когда на наш FastAPI-сервис приходит запрос на скоринг, он за 1 миллисекунду забирает из Redis самую свежую фичу tx_sum_10m и передает ее в оптимизированную ONNX-модель.
# Концептуальный пример логики Spark Streaming
streaming_df = spark.readStream \
    .format("kafka") \
    .option("kafka.bootstrap.servers", "kafka-broker:9092") \
    .option("subscribe", "client_transactions") \
    .load()

# Расчет признака с учетом водяного знака
features_df = streaming_df \
    .withWatermark("event_time", "2 minutes") \
    .groupBy(
        window("event_time", "10 minutes"),
        "client_id"
    ) \
    .sum("amount")

# Запись результата в Redis (Online Feature Store)
features_df.writeStream \
    .foreachBatch(write_to_redis_function) \
    .start()

Потоковая обработка закрывает потребность модели в сверхбыстрых признаках. Однако ML-системам также требуется регулярное переобучение на терабайтах исторических данных, сложные SQL-джоины в хранилищах (Offline Store) и запуск тяжелых проверок качества. Управлять такими многоступенчатыми зависимостями через потоковые скрипты невозможно — для этого требуется специализированный оркестратор пакетных процессов.

Оркестрация пайплайнов данных в Apache Airflow

Оркестрация пайплайнов данных в Apache Airflow

Мы научились обрабатывать потоковые данные в реальном времени с помощью Kafka и Spark Streaming, доставляя свежие признаки в Online Feature Store за миллисекунды. Но ML-система не может жить только в реальном времени. Исторические данные для Offline Feature Store нужно агрегировать терабайтами. Модели деградируют, и их нужно регулярно переобучать на новых данных (Continuous Training). Для этих тяжелых, многошаговых и зависимых друг от друга пакетных (batch) процессов нужен дирижер. В индустрии стандартом де-факто для этой задачи стал Apache Airflow.

Airflow — это не движок для вычислений. Он не перемалывает данные сам, как это делает Spark. Airflow — это планировщик и оркестратор. Он знает, когда запустить задачу, где ее выполнить, что делать в случае ошибки и кто от кого зависит.

Анатомия графа вычислений: DAG

Фундаментальная концепция Airflow — это DAG (Directed Acyclic Graph, направленный ациклический граф).

Пайплайн описывается как набор задач (узлов) и зависимостей между ними (ребер). Слово «направленный» означает, что у процесса есть четкий поток от начала к концу: задача AA выполняется строго до задачи BB. Слово «ациклический» гарантирует, что граф не содержит бесконечных циклов: задача не может зависеть сама от себя ни напрямую, ни через цепочку других задач.

DAG — это логическая схема пайплайна, написанная на Python, которая определяет порядок выполнения задач, их расписание и правила повторного запуска при сбоях.

В коде зависимости задаются интуитивно понятным синтаксисом с помощью битовых операторов сдвига, которые в контексте Airflow переопределены для указания направления:

# Задача extract_data выполнится первой.
# После ее успешного завершения параллельно запустятся train_model и validate_data.
extract_data >> [train_model, validate_data]

Как Airflow выполняет работу: Операторы и Сенсоры

Если DAG — это чертеж, то строительные блоки — это операторы. Оператор описывает конкретное действие, которое должно быть выполнено в рамках одной задачи.

Airflow предлагает сотни готовых операторов для интеграции с любой инфраструктурой:

  • PythonOperator — выполняет обычную Python-функцию (например, легкую трансформацию данных).
  • BashOperator — запускает bash-скрипт.
  • SparkSubmitOperator — отправляет тяжелую задачу на вычисление в Hadoop/Spark кластер.

Особый вид операторов — Сенсоры (Sensors). Это задачи, которые ничего не вычисляют, а просто ждут наступления определенного события. Представьте, что ваш пайплайн переобучения модели кредитного скоринга должен запускаться первого числа каждого месяца. Но данные от партнерского банка выгружаются на S3 с задержкой, иногда второго, а иногда третьего числа. Запускать вычисления по жесткому расписанию нельзя — пайплайн упадет из-за нехватки данных. Здесь используется S3KeySensor: он будет спать и периодически проверять (polling) указанный путь в S3. Как только файл с транзакциями появится, сенсор завершится со статусом «успех», и Airflow запустит следующие по графу задачи очистки и обучения.

Главное правило инженерии данных: Идемпотентность

На технических собеседованиях по MLOps и Data Engineering вопрос об идемпотентности звучит почти всегда.

Идемпотентность — это свойство операции давать один и тот же результат при многократном применении.

В математике примером идемпотентной функции является модуль числа: 5=5|-5| = 5, и повторное применение ничего не меняет: 5=5||-5|| = 5.

В контексте Airflow это означает: если ваша задача упала на середине, и вы перезапустили ее (или весь DAG) за тот же период времени, итоговое состояние системы не должно исказиться.

Пример нарушения идемпотентности: Ваш пайплайн ежедневно агрегирует средний чек клиентов и делает INSERT в таблицу Offline Feature Store. DAG за 15 апреля упал из-за таймаута базы данных, но половина строк успела записаться. Вы нажимаете кнопку "Retry" в интерфейсе Airflow. Задача выполняется успешно, но теперь в базе за 15 апреля лежат дубликаты для половины клиентов. Модель, обученная на таких данных, получит искаженное распределение.

Как сделать пайплайн идемпотентным: Вместо INSERT используйте операцию UPSERT (обновление, если ключ существует, вставка — если нет) или паттерн INSERT OVERWRITE (полная перезапись партиции за конкретный день). Независимо от того, сколько раз Airflow перезапустит задачу за 15 апреля, в базе всегда будет лежать ровно один корректный слепок данных за этот день.

Передача данных между задачами: ловушка XCom

Задачи в Airflow изолированы друг от друга. Они могут выполняться на разных физических серверах (воркерах). Возникает вопрос: как передать результат работы одной задачи в другую?

Для этого в Airflow есть механизм XCom (Cross-Communication). Это небольшая key-value таблица в мета-базе самого Airflow (обычно PostgreSQL), куда задачи могут писать (push) и откуда могут читать (pull) значения.

Красный флаг на собеседовании: сказать, что вы передадите датафрейм pandas (пусть даже на пару сотен мегабайт) из задачи подготовки данных в задачу обучения модели через XCom. Мета-база Airflow не предназначена для хранения данных. Передача больших объектов через XCom приведет к падению базы (OOM) и отказу всего кластера Airflow.

Правильный архитектурный паттерн: Задачи обмениваются через XCom только метаданными (путями, идентификаторами, метриками).

  1. Задача extract_features выгружает данные из БД, сохраняет их в объектное хранилище (S3) и возвращает путь: s3://bucket/features/2023-10/data.parquet.
  2. Этот строковый путь автоматически попадает в XCom.
  3. Задача train_model читает из XCom строку с путем, скачивает Parquet-файл из S3 к себе на воркер и запускает обучение.

Архитектура Continuous Training пайплайна

Соберем изученные концепции в единый пайплайн автоматического переобучения модели (CT), который мы упоминали в первой главе. Как это выглядит в Airflow:

  1. wait_for_data (Sensor): Ждет появления нового батча размеченных данных за прошедшую неделю.
  2. validate_data (Operator): Проверяет схему данных и отсутствие пропусков. Если данные битые — роняет пайплайн и шлет алерт.
  3. extract_features (SparkSubmitOperator): Делает Point-in-Time Join новых целевых меток с историей из Offline Feature Store. Сохраняет датасет на S3.
  4. train_model (PythonOperator / KubernetesPodOperator): Берет путь к датасету из XCom, обучает модель (например, XGBoost), сохраняет артефакт model.pkl в Model Registry.
  5. evaluate_model (Operator): Сравнивает метрики новой модели с текущей продакшен-версией.
  6. conditional_deploy (BranchPythonOperator): Разветвляет логику. Если новая модель лучше (например, ROC-AUC new>oldnew > old) — направляет поток на задачу деплоя. Если хуже — завершает пайплайн.

Airflow связывает разрозненные скрипты дата-саентистов в надежный, отказоустойчивый инженерный конвейер. Он берет на себя ретраи, логирование, распараллеливание и управление зависимостями.

Однако сам код этого графа (Python-файлы DAG-ов), как и код самих моделей, должен как-то попадать на сервер Airflow, тестироваться и версионироваться. И здесь мы переходим от оркестрации данных к оркестрации кода — процессам CI/CD.

Построение CI/CD пайплайнов для ML-сервисов

Построение CI/CD пайплайнов для ML-сервисов

Представьте ситуацию: вы оптимизировали инференс через ONNX, написали асинхронный API на FastAPI и обернули всё это в Docker-контейнер. Локально система работает идеально. Вы отправляете код в репозиторий, ваш коллега вручную собирает образ, выполняет команду обновления в Kubernetes, и... продакшен падает. Оказалось, что в новой Pydantic-схеме опечатка в типе данных, а переменная окружения для подключения к Redis (нашему Online Feature Store) не передана.

Ручной деплой — это узкое горлышко и гарантия человеческой ошибки. Чтобы ML-система была по-настоящему промышленной, процесс доставки кода и моделей должен быть автоматизирован так же надежно, как конвейер на заводе.

Трехмерное пространство: CI, CD и CT

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

  1. CT (Continuous Training) — автоматическое переобучение модели на свежих данных (мы реализовали это с помощью направленных графов в Apache Airflow).
  2. CI (Continuous Integration) — непрерывная интеграция кода. Проверка того, что новые изменения от разработчиков не ломают текущую логику.
  3. CD (Continuous Deployment) — непрерывное развертывание. Автоматическая доставка проверенного кода и артефактов на серверы (в наш Kubernetes-кластер).

CI/CD — практика непрерывного слияния рабочих копий разработчиков в общую ветку репозитория и автоматизированного развертывания результата на целевых средах без ручного вмешательства.

Если Airflow отвечает за то, чтобы модель оставалась актуальной математически, то CI/CD отвечает за то, чтобы сервис, оборачивающий эту модель, работал стабильно с инженерной точки зрения.

Continuous Integration: Проверка на прочность

Процесс CI запускается автоматически при каждом коммите (например, git push) в репозиторий. Для ML-сервиса этот этап должен ответить на главный вопрос: безопасно ли это изменение?

Типичный CI-пайплайн состоит из серии последовательных проверок:

  1. Линтеры (Linters). Программы, которые анализируют код на соответствие стандартам оформления и ищут синтаксические ошибки без его выполнения (например, flake8 или ruff). Они гарантируют, что весь код в команде выглядит единообразно.
  2. Тайп-чекеры (Type-checkers). В Python типы динамические, что часто приводит к ошибкам в рантайме. Инструменты вроде mypy статически проверяют аннотации типов. Если функция ожидает int, а вы передаете str, тайп-чекер остановит пайплайн еще до запуска тестов.
  3. Unit-тесты. Проверка изолированных кусков логики (например, через pytest).

Специфика CI для машинного обучения

Главное правило MLOps: мы никогда не обучаем модели внутри CI-пайплайна. Обучение требует GPU, часов времени и огромных объемов данных. CI должен отрабатывать за 2–5 минут.

Вместо обучения, в CI мы тестируем интеграцию модели и кода:

  • Тест загрузки: удается ли успешно десериализовать веса (например, загрузить ONNX-граф) за отведенное время?
  • Тест контрактов: совпадает ли размерность входного тензора с тем, что генерирует наш код подготовки данных?
  • Тест инференса: если подать на вход модели фиксированный вектор-заглушку, вернет ли она ожидаемый результат без падения по памяти (OOM)?

Continuous Deployment: Упаковка и доставка

Если CI-пайплайн загорелся зеленым (все тесты пройдены), запускается этап CD. Его задача — превратить исходный код в готовый к запуску артефакт и доставить его пользователям.

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

  1. Сборка образа. Пайплайн выполняет команду docker build. При этом используются механизмы кэширования слоев, чтобы не скачивать тяжелые ML-библиотеки заново при каждом деплое.
  2. Публикация (Push). Готовый образ помечается уникальным тегом (обычно это хэш git-коммита, например credit-scoring:a1b2c3d) и отправляется в Container Registry — защищенное хранилище Docker-образов.

Но как этот новый образ попадет в Kubernetes? Исторически CI-сервер сам подключался к кластеру и выполнял команду обновления. Этот подход называется Push-моделью, и сегодня он считается устаревшим из-за проблем с безопасностью (CI-серверу нужны админские права от продакшена).

GitOps: Смена парадигмы

Современный стандарт доставки в Kubernetes — это подход GitOps.

GitOps — парадигма управления инфраструктурой, при которой Git-репозиторий выступает единственным источником истины для декларативного описания системы.

Вместо того чтобы «проталкивать» изменения снаружи, мы используем Pull-модель. Внутри самого кластера Kubernetes устанавливается специальный агент (самые популярные решения — ArgoCD или Flux).

Агент непрерывно наблюдает за отдельным Git-репозиторием, где хранятся конфигурационные файлы (манифесты Kubernetes). Как только агент замечает, что в репозитории изменилась версия образа (например, с v1.0 на v1.1), он самостоятельно скачивает новый образ из Container Registry и инициирует стратегию бесшовного обновления (Rolling Update, которую мы настраивали для наших Подов).

Жизненный цикл одного изменения

Давайте сведем все воедино на примере нашего сервиса кредитного скоринга. Разработчик добавил новую фичу — проверку возраста клиента перед отправкой запроса в модель.

Этап Инструмент Что происходит под капотом
Код Git Разработчик делает git push в ветку feature/age-check.
CI GitLab CI / GitHub Actions Запускается mypy (проверяет типы), затем pytest (проверяет логику отсева по возрасту). Тесты пройдены.
CD (Build) Docker Собирается новый образ scoring-api:v2.0 и отправляется в Container Registry.
CD (Manifest) Скрипт в пайплайне Пайплайн автоматически делает коммит в инфраструктурный репозиторий, меняя в манифесте Deployment версию образа на v2.0.
GitOps ArgoCD Агент внутри K8s видит новый коммит, сверяет желаемое состояние с текущим и плавно заменяет старые Поды на новые.

В этой архитектуре разработчику не нужен доступ к продакшен-серверам. Вся система управляется исключительно через коммиты в Git. Если новая версия содержит баг, откат (Rollback) сводится к нажатию кнопки "Revert" в истории Git — и ArgoCD мгновенно вернет кластер к предыдущему стабильному состоянию.

Теперь наша система автоматизирована от момента написания кода до его запуска в кластере. Но как только модель начинает принимать реальный трафик, возникает новая задача: как понять, что поведение пользователей изменилось и модель начала ошибаться, даже если сам код работает идеально?

Мониторинг моделей и детектирование сдвига данных

Мониторинг моделей и детектирование сдвига данных

Код успешно прошел линтеры, Docker-образ собран, а агент ArgoCD развернул новые поды в Kubernetes. Графики показывают идеальную картину: потребление CPU стабильно, задержка ответов (Latency) не превышает 15 мс, HTTP-ошибок нет. С точки зрения классической инфраструктуры сервис работает безупречно. Но спустя три месяца бизнес сообщает о резком росте дефолтов по кредитам, одобренным нашей моделью. Технически система функционировала идеально, но статистически она потерпела крах.

В машинном обучении успешный HTTP-ответ 200 OK не гарантирует корректности бизнес-логики. Модель может бесперебойно выдавать абсолютно неверные предсказания, если изменился мир, в котором она работает.

Проблема отложенной истины

В идеальном мире мы бы мониторили качество модели напрямую: сравнивали предсказание с реальностью в режиме реального времени и строили график ROC-AUC или F1-score.

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

Ground Truth (Истинная метка) — фактический результат события, который модель пыталась предсказать.

В задаче кредитного скоринга малого бизнеса (КМБ) мы узнаем, вернул ли клиент кредит, только через 6–12 месяцев. Если мы будем ждать Ground Truth для оценки качества модели, бизнес понесет колоссальные убытки до того, как мы заметим деградацию.

Раз мы не можем измерить точность напрямую, нам нужны прокси-метрики. Мы переносим фокус с оценки результата на оценку входных данных и распределения ответов. Если статистический профиль входящих запросов совпадает с тем, на котором модель обучалась, мы предполагаем, что она работает корректно.

Метрики: Системные против ML

Мониторинг ML-сервиса требует двух независимых контуров наблюдения.

Категория Что измеряем Инструменты Сигнал об ошибке
Системные метрики Утилизация ресурсов (CPU, RAM), Latency, Throughput, коды ответов. Prometheus, Grafana, Kubernetes (OOMKilled). Поды падают, сервис отвечает 503 Service Unavailable, таймауты.
ML-метрики Распределение признаков (Data Drift), распределение предсказаний, доля пропущенных значений. Evidently, Spark Streaming, кастомные экспортеры. Данные сместились, модель начала одобрять 90% заявок вместо исторических 15%.

Математика детектирования сдвига: PSI

Чтобы автоматизировать поиск аномалий в данных, нам нужен математический аппарат, способный сравнить два распределения: эталонное (из Offline Feature Store на момент обучения) и текущее (из потока запросов к API).

В банковском секторе стандартом де-факто является индекс стабильности популяции.

Population Stability Index (PSI) — метрика, показывающая, насколько сильно изменилось распределение переменной по сравнению с эталонным состоянием.

Формула расчета PSI для одного признака, разбитого на интервалы (корзины):

PSI=i=1n(ActualiExpectedi)×ln(ActualiExpectedi)PSI = \sum_{i=1}^{n} (Actual_i - Expected_i) \times \ln\left(\frac{Actual_i}{Expected_i}\right)

Где:

  • nn — количество интервалов, на которые мы разбили значения признака.
  • ActualiActual_i — доля наблюдений, попавших в ii-й интервал в текущем трафике (например, за последние 24 часа).
  • ExpectediExpected_i — доля наблюдений в этом же интервале в обучающей выборке.
  • ln\ln — натуральный логарифм.

Практический пример: Мы мониторим признак monthly_revenue (ежемесячная выручка бизнеса). При обучении модели мы разбили выручку на корзины. Возьмем одну корзину: от 100 тыс. до 500 тыс. руб. В обучающей выборке в эту корзину попадало 20% клиентов (Expected=0.20Expected = 0.20). Вчера маркетинг запустил новую рекламную кампанию, и теперь в эту корзину попадает 40% клиентов (Actual=0.40Actual = 0.40).

Считаем вклад этой корзины в общий PSI: Разница долей: 0.400.20=0.200.40 - 0.20 = 0.20 Отношение долей: 0.40/0.20=20.40 / 0.20 = 2 Логарифм отношения: ln(2)0.693\ln(2) \approx 0.693 Вклад корзины: 0.20×0.6930.1380.20 \times 0.693 \approx 0.138

Просуммировав такие значения по всем корзинам, мы получаем итоговый PSI признака. Интерпретация результата строго стандартизирована:

  • PSI<0.1PSI < 0.1 — изменений нет, зеленая зона.
  • 0.1PSI<0.20.1 \leq PSI < 0.2 — есть небольшие изменения, желтая зона (требует внимания аналитика).
  • PSI0.2PSI \geq 0.2 — значительный сдвиг данных, красная зона (модель больше нельзя использовать без переобучения).

Архитектура системы мониторинга

Вычисление PSI и других статистик — это ресурсоемкая операция. Если мы добавим этот расчет прямо в код FastAPI-приложения, обрабатывающего инференс, мы заблокируем Event Loop и обрушим пропускную способность сервиса. Мониторинг должен быть асинхронным и изолированным.

Архитектурный паттерн выглядит следующим образом:

  1. Асинхронное логирование (Shadow Logging). FastAPI-сервис при ответе на запрос не занимается аналитикой. Он берет вектор признаков, само предсказание, добавляет метку времени и отправляет этот JSON-документ в топик Kafka (например, model_inference_logs). Отправка происходит в режиме "fire-and-forget", не замедляя ответ клиенту.
  2. Агрегация в потоке. Отдельное приложение на базе Spark Structured Streaming читает топик Kafka. Оно накапливает микро-батчи (например, за окно в 1 час), вычисляет распределение признаков (ActualActual) и сравнивает их с закэшированными эталонными значениями (ExpectedExpected) из Offline Store.
  3. Экспорт метрик. Spark-приложение выставляет рассчитанные значения PSI по каждому признаку в формате, понятном для систем мониторинга (Prometheus).
  4. Визуализация и алертинг. Prometheus скрейпит (собирает) эти метрики. Grafana отрисовывает тепловую карту: по оси X — время, по оси Y — признаки, цвет ячейки — значение PSI.

Замыкание цикла: от мониторинга к Continuous Training

Мониторинг ради красивых дашбордов бесполезен. Истинная ценность MLOps заключается в автоматической реакции на деградацию.

Когда Prometheus фиксирует, что PSI ключевого признака (или самого предсказания) превысил порог 0.20.2 и держится на этом уровне более суток, срабатывает Alertmanager.

Вместо того чтобы просто разбудить инженера посреди ночи, Alertmanager отправляет HTTP-webhook в Apache Airflow. Этот триггер запускает DAG процесса Continuous Training (CT). Airflow инициирует сбор свежего датасета (с учетом новых распределений), переобучает модель, прогоняет обновленные веса через CI-пайплайн и, если новые метрики качества удовлетворяют бизнес-требованиям, инициирует деплой новой версии.

Таким образом, мониторинг становится сенсором, который делает систему машинного обучения самовосстанавливающейся.

Мы рассмотрели мониторинг классических табличных моделей, где распределение признаков поддается строгой математической оценке. Но как быть, если наша модель генерирует не число вероятности дефолта, а связный текст, и работает не с таблицами, а с неструктурированными промптами? Об этом пойдет речь при проектировании архитектуры генеративных систем.

Архитектура LLM и RAG-систем в продакшене

Архитектура LLM и RAG-систем в продакшене

Классическая ML-модель выдает конкретное число или класс — вероятность дефолта 0.850.85 или категорию «Фрод». Мы научились упаковывать такие модели в контейнеры, балансировать нагрузку и отслеживать сдвиг распределений. Но что происходит, когда ответом системы становится абзац связного текста длиной в 500 слов, который генерируется токен за токеном? С приходом больших языковых моделей (LLM) классический стек MLOps не ломается, но требует серьезного расширения: от способов работы с памятью видеокарт до совершенно новых подходов к мониторингу качества.

Узкое горлышко инференса LLM: память и контекст

В классическом инференсе главная задача — быстрее перемножить матрицы. В LLM архитектуре (Transformer) возникает новая проблема: каждое следующее слово генерируется на основе всех предыдущих. Чтобы не пересчитывать весь текст с нуля ради одного нового слова, системы сохраняют промежуточные вычисления в оперативной памяти GPU. Этот механизм называется KV Cache (Key-Value Cache).

Проблема в том, что длина запросов пользователей непредсказуема. Один клиент пишет два слова, другой вставляет лог ошибки на три страницы. Если выделять память под KV Cache статично (с запасом под максимальную длину), память GPU моментально фрагментируется и заканчивается. Модель простаивает, ожидая свободной памяти, и мы получаем ошибку OOM (Out Of Memory), даже если вычислительные ядра загружены всего на 20%.

Инженерное решение пришло из классических операционных систем — страничная память. Алгоритм PagedAttention (реализованный, например, в библиотеке vLLM) разбивает KV Cache на блоки фиксированного размера (страницы). Память выделяется динамически по мере генерации текста, устраняя фрагментацию. Это позволяет обрабатывать в 3–4 раза больше одновременных запросов на том же железе, применяя уже знакомый нам динамический батчинг на уровне отдельных токенов, а не целых запросов.

Проблема изоляции: почему LLM нужна помощь

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

Переобучать или дообучать (Fine-tuning) гигантскую модель каждый день при обновлении регламентов — астрономически дорого и долго. Вместо этого индустрия использует паттерн RAG (Retrieval-Augmented Generation — генерация, дополненная поиском).

Идея RAG проста: прежде чем отправить вопрос пользователя в LLM, система идет во внутреннюю базу знаний, находит там релевантные документы и подкладывает их в запрос как шпаргалку.

Архитектура RAG делится на два независимых пайплайна:

  1. Offline-пайплайн (Подготовка данных): Внутренние документы (PDF, страницы Confluence) разбиваются на небольшие куски текста — чанки (chunks). Специальная ML-модель переводит каждый чанк в эмбеддинг — числовой вектор, отражающий смысл текста. Эти векторы сохраняются в векторную базу данных.
  2. Online-пайплайн (Инференс): Пользователь задает вопрос. Тот же алгоритм превращает вопрос в вектор. База данных ищет векторы документов, наиболее близкие к вектору вопроса. Найденный текст вклеивается в промпт, и LLM формирует финальный ответ.

Векторный поиск и математика близости

Как база данных понимает, что текст «Условия выдачи ипотеки» релевантен вопросу «Как купить квартиру в кредит»? Векторы этих фраз в многомерном пространстве будут находиться рядом.

Для измерения расстояния между векторами чаще всего используется косинусное сходство:

Cosine Similarity=ABAB\text{Cosine Similarity} = \frac{\mathbf{A} \cdot \mathbf{B}}{\|\mathbf{A}\| \|\mathbf{B}\|}

Где A\mathbf{A} — это вектор вопроса пользователя, B\mathbf{B} — вектор фрагмента документа из базы, AB\mathbf{A} \cdot \mathbf{B} — их скалярное произведение, а знаменатель — произведение их длин.

Результат формулы лежит в диапазоне от 1-1 до 11. Значение, близкое к 11, означает, что векторы указывают в одном направлении (смысл совпадает). Значение около 00 говорит об отсутствии связи.

Например, если вектор запроса A\mathbf{A} («лимиты по кредиткам») и вектор документа B\mathbf{B} («максимальная сумма кредитной карты») имеют косинусное сходство 0.920.92, этот документ будет извлечен первым и передан в LLM.

Для хранения и быстрого поиска по миллионам таких векторов используются векторные базы данных (Milvus, Qdrant или расширение pgvector для PostgreSQL). Они строят специальные индексы (например, HNSW), которые позволяют находить ближайших соседей за миллисекунды без необходимости сравнивать запрос с каждым документом в базе.

От RAG к Агентам: модели, которые действуют

RAG решает проблему чтения данных, но система остается пассивной. Что если пользователь просит: «Заблокируй мою карту, я ее потерял»? Поиском по документам тут не обойтись — нужно совершить действие в другой системе.

Здесь на сцену выходят AI-агенты. Агент — это архитектурный паттерн, где LLM выступает в роли «мозга», которому выдали набор инструментов (Tools).

Вместо того чтобы просто генерировать текст, модель обучена возвращать структурированный JSON с названием функции и аргументами.

  1. Пользователь пишет: «Узнай статус заявки №456».
  2. LLM понимает, что для ответа ей нужен внешний инструмент, и генерирует команду: call: get_application_status(id=456).
  3. Бэкенд перехватывает эту команду, делает реальный HTTP-запрос к CRM-системе и получает ответ: «На рассмотрении».
  4. Бэкенд возвращает этот статус обратно в LLM, и та формулирует человечный ответ: «Ваша заявка сейчас находится на рассмотрении».

Агенты превращают LLM из умного справочника в полноценный интерфейс управления инфраструктурой.

Мониторинг текста: LLM-as-a-Judge

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

Для мониторинга RAG-систем в продакшене используется подход LLM-as-a-Judge (LLM в роли судьи). Мы берем отдельную, часто более дешевую или специализированную модель, и поручаем ей оценивать логи работы основной системы по трем осям (RAG Triad):

  1. Context Relevance (Релевантность контекста): Нашел ли векторный поиск действительно полезные документы для ответа на вопрос?
  2. Groundedness (Обоснованность): Опирается ли финальный ответ LLM только на найденные документы, или она приплела факты из интернета?
  3. Answer Relevance (Релевантность ответа): Отвечает ли финальный текст на изначальный вопрос пользователя?

Оценки (обычно от 1 до 5) логируются асинхронно и выводятся на дашборды. Если метрика Groundedness начинает падать, это сигнал для инженера: модель начала фантазировать. Возможно, изменилась структура документов, векторный поиск стал приносить мусор, и LLM пытается выкрутиться, генерируя ответы из своих внутренних весов.

Мы прошли путь от упаковки простой модели в Docker до оркестрации распределенных систем, потоковой обработки признаков и архитектуры агентов. Все эти компоненты — кирпичики, из которых строятся реальные продукты. Настало время собрать их воедино и посмотреть, как выстраивать архитектуру всей ML-системы целиком при проектировании с нуля.

Проектирование ML-систем на собеседовании

Проектирование ML-систем на собеседовании

Можно безупречно написать Dockerfile, настроить автомасштабирование в Kubernetes и ускорить инференс квантованием, но провалить финальное собеседование по одной причине. Интервьюер просит: «Спроектируй систему кредитного скоринга для малого бизнеса», а кандидат начинает хаотично рисовать на доске квадратики с логотипами Kafka, FastAPI и Redis, не связав их единой логикой данных и бизнес-целью. ML System Design — это не проверка знания инструментов. Это проверка вашей способности собрать из них работающий, отказоустойчивый и экономически целесообразный конвейер.

В этой финальной главе мы объединим все компоненты, изученные ранее, в единую архитектуру. Мы пройдем по классическому фреймворку системного дизайна на примере реальной задачи из сегмента КМБ (кредитование малого бизнеса).

Фреймворк ответа на ML System Design

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

  1. Формирование требований (Requirements). Перевод бизнес-задачи в инженерные ограничения.
  2. Потоки данных (Data Flow). Как данные собираются для обучения и как доставляются в момент инференса.
  3. Модель и вычисления (Model & Inference). Как модель упакована и как она отдает предсказания.
  4. Инфраструктура и мониторинг (Ops & Monitoring). Как система обновляется и как мы узнаем, что она сломалась.

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

Задача от интервьюера: Нам нужна система автоматического принятия решений по кредитам для малого бизнеса. Клиент нажимает кнопку в приложении, и мы должны либо одобрить лимит, либо отправить заявку на ручное рассмотрение.

Шаг 1. Требования и ограничения

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

Категория Вопрос интервьюеру Пример ответа (ограничение системы)
Бизнес-цель Что важнее: не выдать плохой кредит (Precision) или одобрить максимум хороших (Recall)? Важнее Precision, цена дефолта бизнеса высока.
Latency Какое максимальное время ожидания ответа для пользователя? Не более 500500 мс, иначе клиент уходит.
Throughput Какая пиковая нагрузка ожидается? До 100100 запросов в секунду (RPS).
Данные Есть ли неструктурированные данные (например, сканы налоговых деклараций)? Да, клиенты могут прикреплять PDF-отчеты.

Зафиксировав 500500 мс и 100100 RPS, вы уже отсекли варианты с тяжелой пакетной обработкой (Batch) в пользу синхронного API (Real-time), а наличие PDF-документов потребует интеграции LLM-агентов.

Шаг 2. Архитектура данных (Feature Engineering Pipeline)

Модель не может работать с сырым client_id. Ей нужны признаки (фичи). На этом этапе вы должны показать, как решается проблема Training-Serving Skew.

Для скоринга КМБ нам нужны два типа признаков:

  1. Медленные (Статические): Возраст компании, отрасль, кредитная история за прошлые годы. Они пересчитываются редко.
  2. Быстрые (Динамические): Оборот по расчетному счету за последние 10 минут, количество подозрительных транзакций за день.

Как это нарисовать: Вы проектируете двухконтурный Feature Store. Сырые транзакции летят в брокер сообщений (Kafka). Оттуда движок потоковой обработки (Spark Streaming) агрегирует их в микро-батчах, отсекая опоздавшие события с помощью водяных знаков (Watermarks). Результат записывается в два места:

  • В Offline Store (например, S3 или Hadoop) — для будущего переобучения моделей. Здесь применяется механизм Point-in-Time Join, чтобы при сборке датасета не заглянуть в будущее.
  • В Online Store (Redis) — для мгновенного доступа при инференсе. Redis обеспечит чтение признака по ключу client_id за 11 мс, что критично для нашего бюджета в 500500 мс.

Шаг 3. Инференс и вычисления

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

Запрос от клиента поступает в Kubernetes-кластер. Балансировщик направляет его в Pod с нашим приложением на FastAPI.

Логика внутри FastAPI (Обогащение):

  1. Приложение валидирует входящий JSON с помощью Pydantic-схемы.
  2. Делает асинхронный вызов в Redis, забирая вектор признаков.
  3. Если в запросе есть текстовая выписка (PDF), FastAPI передает ее AI-Агенту. Агент использует LLM для извлечения ключевых сущностей (например, "чистая прибыль"). Чтобы LLM не съела всю память при длинных документах, мы разворачиваем ее с использованием механизма PagedAttention.

Оптимизация модели: Классическая модель скоринга (например, градиентный бустинг) не запускается в нативном Python-окружении. На этапе CD-пайплайна мы экспортировали ее в статический граф ONNX и применили квантование, снизив точность весов до INT8. Внутри FastAPI инференс работает через ONNX Runtime. Поскольку матричные вычисления — это CPU-bound задача, мы запускаем предсказание в отдельном пуле потоков (или процессов), чтобы не заблокировать Event Loop асинхронного веб-сервера.

Если RPS резко вырастет со 100100 до 500500, сработает Horizontal Pod Autoscaler (HPA) в Kubernetes, который отследит утилизацию CPU и поднимет дополнительные реплики FastAPI-сервиса.

Шаг 4. Мониторинг и Continuous Training (CT)

Выдав предсказание клиенту, система не заканчивает работу. Одобренные кредиты могут стать дефолтными через полгода — это проблема задержки Ground Truth. Мы не можем ждать полгода, чтобы понять, что модель сломалась.

Замыкаем цикл:

  1. FastAPI-сервис асинхронно (чтобы не тормозить ответ клиенту) отправляет входные признаки и выданный скор-балл в отдельный топик Kafka (например, model_logs).
  2. Ежедневно Apache Airflow запускает DAG, который вычитывает эти логи и сравнивает распределение сегодняшних признаков с распределением, на котором модель обучалась.
  3. Airflow рассчитывает метрику Population Stability Index (PSI).
  4. Если PSI >0.2> 0.2 (произошел Data Drift — например, пошли клиенты из новой отрасли), Alertmanager отправляет уведомление команде и автоматически триггерит следующий DAG в Airflow — пайплайн Continuous Training.
  5. Выполняется переобучение на свежих данных из Offline Feature Store.
  6. Новая модель проходит CI-проверки контрактов, упаковывается в Docker-образ вместе с зависимостями (используя паттерн Bake-in для весов) и через GitOps-агента (ArgoCD) раскатывается в кластер с помощью стратегии Rolling Update.

Компромиссы (Trade-offs) в архитектуре

На собеседовании вас обязательно попытаются сбить с толку, предложив изменить условия. Умение аргументированно выбирать между двумя решениями — главный маркер Senior/Middle+ инженера.

Ситуация 1: "А давайте считать скоринг не в реальном времени, а заранее?" Batch-инференс против Real-time. Вы можете предложить рассчитывать скор-баллы для всех существующих клиентов банка каждую ночь (Airflow + Spark) и просто складывать готовые ответы в Redis. Плюсы: Нулевая задержка при нажатии кнопки (просто чтение из базы), не нужно держать FastAPI под высокой нагрузкой. Минусы: Не работает для новых клиентов (у которых нет истории в банке). Не учитывает транзакции, совершенные 5 минут назад.

Ситуация 2: "Модель слишком тяжелая, мы не укладываемся в 500500 мс. Что делать?" Вы должны предложить каскадную архитектуру. Сначала запускается легкая, быстрая модель (например, логистическая регрессия на базовых фичах). Если она дает однозначный ответ (очень хороший или очень плохой клиент), мы сразу возвращаем результат. Если случай пограничный — запрос отправляется в тяжелую модель (LLM или ансамбль бустингов), для которой клиент видит плашку "Анализируем ваши документы...", оправдывающую долгое ожидание.

Сборка картины воедино

Архитектура ML-системы — это перевод математики в инженерию. Мы начали с проблемы зависимостей и изоляции (Docker), научились управлять этими контейнерами при нагрузке (Kubernetes) и написали для них правильный асинхронный интерфейс (FastAPI). Мы ускорили саму математику (ONNX, квантование) и обеспечили ее свежими данными без утечек (Feature Store, Kafka). Наконец, мы автоматизировали доставку кода (GitOps) и переобучение при деградации (Airflow, PSI).

Каждый инструмент в вакансии MLOps-инженера решает строго одну фундаментальную проблему жизненного цикла модели. Понимая, на каком этапе находится узкое горлышко, вы никогда не ошибетесь с выбором технологии на собеседовании и в продакшене.