Apache Airflow для ML-пайплайнов: настройка, оркестрация, автоматизация

Проектируем и внедряем системы искусственного интеллекта: от прототипа до production-ready решения. Наша команда объединяет экспертизу в машинном обучении, дата-инжиниринге и MLOps, чтобы AI работал не в лаборатории, а в реальном бизнесе.
Показано 1 из 1Все 1564 услуг
Apache Airflow для ML-пайплайнов: настройка, оркестрация, автоматизация
Средний
~3-5 дней
Часто задаваемые вопросы

Направления AI-разработки

Этапы разработки AI-решения

Последние работы

  • image_website-b2b-advance_0.webp
    Разработка сайта компании B2B ADVANCE
    1361
  • image_web-applications_feedme_466_0.webp
    Разработка веб-приложения для компании FEEDME
    1251
  • image_websites_belfingroup_462_0.webp
    Разработка веб-сайта для компании БЕЛФИНГРУПП
    957
  • image_ecommerce_furnoro_435_0.webp
    Разработка интернет магазина для компании FURNORO
    1189
  • image_logo-advance_0.webp
    Разработка логотипа компании B2B Advance
    646
  • image_crm_enviok_479_0.webp
    Разработка веб-приложения для компании Enviok
    929

При оркестрации ML-пайплайнов в продакшене часто возникает проблема: нужно связать препроцессинг на CPU-нодах, обучение на GPU-нодах с разными конфигурациями, валидацию качества по метрикам и автоматический деплой — и всё это по расписанию, с откатами при падении метрик. Представьте: ежедневное переобучение модели fraud detection, где загрузка данных из S3, препроцессинг на 4 CPU, обучение на 1 GPU, валидация F1 и деплой в staging требуют координации. Без оркестрации инженер вручную запускает скрипты, следит за логами и при сбое теряет часы. Apache Airflow автоматизирует этот процесс через DAG-графы, KubernetesExecutor для динамического выделения ресурсов и интеграцию с MLflow. Мы в течение 10+ лет настраиваем Airflow для ML-пайплайнов — от небольших команд до enterprise-кластеров с 500+ DAGов. Наш опыт включает проекты с fraud detection, NLP, Computer Vision, где автоматизация пайплайнов сократила время на эксперименты на 40-60% и снизила число инцидентов при деплое в 3 раза. По сравнению с ручным запуском, Airflow снижает время on-бординга новых моделей в 2-3 раза. Экономия на DevOps-часах достигает 30-50%. Мы гарантируем SLA 99.9% и предоставляем документацию, мониторинг и обучение команды.

Как Apache Airflow решает проблемы ML-оркестрации?

Airflow решает ключевые проблемы оркестрации ML: гетерогенность ресурсов (CPU/GPU), управление зависимостями между задачами, повторяемость и отказоустойчивость. Каждый шаг пайплайна — отдельная задача в DAG: подготовка данных на стандартном поде, обучение на GPU-поде с tolerations, валидация качества через Python-оператор и промоция модели. При падении качества (F1 < 0.90) DAG останавливается с ошибкой, что предотвращает выкат плохой модели. Все метрики логируются в MLflow, что позволяет сравнивать эксперименты. Airflow с KubernetesExecutor лучше CeleryExecutor для ML-задач в 2 раза по изоляции ресурсов: каждый под с GPU изолирован, не влияет на соседние задачи. Это критично при смешанных рабочих нагрузках.

Сравнение исполнителей Airflow для ML

Исполнитель Изоляция ресурсов Поддержка GPU Сложность Сценарий использования
KubernetesExecutor Полная (каждая задача в своём поде) Да Средняя ML-пайплайны с GPU, гибридные кластеры
CeleryExecutor Нет (задачи на общих воркерах) Ограниченная Низкая ETL, небольшие ML-задачи без GPU
LocalExecutor Нет Нет Минимальная Разработка, тестирование

В чем разница между Airflow и Kubeflow для ML?

Аспект Airflow Kubeflow Pipelines
Тип задач Универсальный оркестратор (ETL + ML) Только ML-пайплайны
Примитивы DAG, операторы, сенсоры Components, pipelines, metrics
Интеграция Любые системы (S3, BigQuery, MLflow) Нативная интеграция с K8s и Kubeflow
Когда выбрать Уже есть Airflow, нужна гибкость ML-центричная команда, только K8s

Airflow выигрывает в универсальности, Kubeflow — в глубине ML-интеграции. Если команда уже использует Airflow для ETL, миграция ML-пайплайнов на него сокращает затраты на инфраструктуру на 30%.

Установка с KubernetesExecutor

# Установка через Helm (рекомендуется)
helm repo add apache-airflow https://airflow.apache.org
helm upgrade --install airflow apache-airflow/airflow \
  --namespace airflow \
  --create-namespace \
  --set executor=KubernetesExecutor \
  --set config.logging.logging_level=INFO \
  --values airflow-values.yaml

ML-пайплайн как Airflow DAG

from airflow import DAG
from airflow.providers.cncf.kubernetes.operators.pod import KubernetesPodOperator
from airflow.operators.python import PythonOperator
from airflow.operators.trigger_dagrun import TriggerDagRunOperator
from datetime import datetime, timedelta

default_args = {
    "owner": "ml-team",
    "retries": 2,
    "retry_delay": timedelta(minutes=5),
    "on_failure_callback": notify_on_slack,
}

with DAG(
    "fraud_detection_training",
    default_args=default_args,
    schedule="0 2 * * 1",  # по понедельникам в 2:00
    start_date=datetime(2025, 1, 1),
    catchup=False,
    tags=["ml", "fraud-detection"],
) as dag:

    # Подготовка данных — на обычном поде
    prepare_data = KubernetesPodOperator(
        task_id="prepare_data",
        image="ml-pipeline:latest",
        cmds=["python", "prepare_data.py"],
        arguments=["--date={{ ds }}", "--output=s3://bucket/features/{{ ds }}/"],
        namespace="ml-pipelines",
        resources={"request_memory": "4Gi", "request_cpu": "2"},
        get_logs=True,
        is_delete_operator_pod=True,
    )

    # Обучение — на GPU поде
    train_model = KubernetesPodOperator(
        task_id="train_model",
        image="ml-pipeline-gpu:latest",
        cmds=["python", "train.py"],
        arguments=[
            "--data=s3://bucket/features/{{ ds }}/",
            "--run-name=fraud-{{ ds }}",
        ],
        namespace="ml-pipelines",
        resources={
            "request_memory": "32Gi",
            "request_cpu": "8",
            "limit_gpu": "1",
        },
        annotations={"nvidia.com/gpu": "1"},
        tolerations=[{"key": "nvidia.com/gpu", "operator": "Exists", "effect": "NoSchedule"}],
        get_logs=True,
    )

    # Evaluation gate — Python оператор (дешево)
    def check_model_quality(**context):
        import mlflow
        client = mlflow.tracking.MlflowClient()
        run = client.search_runs(
            experiment_ids=[EXPERIMENT_ID],
            filter_string=f"tags.run_date = '{context['ds']}'",
            order_by=["metrics.f1 DESC"],
            max_results=1
        )[0]
        f1 = run.data.metrics.get("test_f1", 0)
        if f1 < 0.90:
            raise ValueError(f"Model quality too low: F1={f1:.3f} < 0.90")
        context["ti"].xcom_push(key="run_id", value=run.info.run_id)

    quality_gate = PythonOperator(
        task_id="quality_gate",
        python_callable=check_model_quality,
    )

    # Промоция — только если quality_gate прошёл
    promote_model = KubernetesPodOperator(
        task_id="promote_to_staging",
        image="ml-pipeline:latest",
        cmds=["python", "promote_model.py"],
        arguments=["--run-id={{ ti.xcom_pull(task_ids='quality_gate', key='run_id') }}"],
        namespace="ml-pipelines",
    )

    # Зависимости
    prepare_data >> train_model >> quality_gate >> promote_model

TaskFlow API (современный подход)

from airflow.decorators import dag, task

@dag(schedule="0 2 * * 1", start_date=datetime(2025, 1, 1))
def ml_pipeline():
    @task
    def prepare_data(execution_date: str) -> str:
        # Подготовка данных
        return f"s3://bucket/features/{execution_date}/"

    @task
    def train_model(data_path: str) -> dict:
        # Запуск обучения (или триггер внешнего job)
        return {"run_id": "xxx", "f1": 0.924}

    @task
    def promote_if_good(metrics: dict) -> None:
        if metrics["f1"] >= 0.90:
            promote_to_staging(metrics["run_id"])

    data = prepare_data()
    metrics = train_model(data)
    promote_if_good(metrics)

ml_pipeline()

Мониторинг Airflow DAG

Airflow UI показывает: статус каждого запуска, длительность каждого task, логи. Интеграция с Prometheus через airflow-exporter: airflow_dag_run_duration_seconds, airflow_task_fail_count. Алерт при failed task через Slack/PagerDuty через on_failure_callback. Для глубокого мониторинга ML-метрик (дрейф данных, распределение предсказаний) рекомендуется интегрировать Evidently AI или WhyLabs — они триггерят повторное обучение при дрейфе.

Типичные ошибки при настройке Airflow для ML
  • Использование CeleryExecutor с GPU-задачами — приводит к конфликтам памяти.
  • Отсутствие retry для препроцессинга — при кратковременных сбоях S3 пайплайн падает.
  • Игнорирование timeouts для долгих задач обучения — DAG зависает навсегда.
  • Неправильные tolerations для GPU-нод — поды не попадают на GPU-кластер.

Чтобы избежать этого, мы используем KubernetesExecutor, задаём явные таймауты и тестируем пайплайн на staging.

Что входит в настройку Airflow под ключ

Мы предоставляем полный цикл настройки: аудит текущей инфраструктуры, проектирование архитектуры DAG с учётом ML-специфики (GPU, большие данные), установка и конфигурация Airflow на Kubernetes с Helm, настройка мониторинга (Prometheus + Grafana) и алертинга, интеграция с MLflow, написание 5-10 кастомных DAG под ваши задачи, обучение команды и техническая поддержка на этапе эксплуатации.

Свяжитесь с нами для бесплатной консультации — мы проанализируем ваш проект и предложим оптимальную архитектуру. Закажите внедрение Airflow — получите стабильный ML-пайплайн за недели, а не месяцы.

MLOps: инфраструктура для обучения, деплоя и мониторинга ML-моделей

Модель обучена, метрики — F1 0.94 на валидации. Через три месяца в продакшене качество падает на 12%. Никто не знает, когда именно — нет мониторинга. Нельзя быстро переобучить — обучающий скрипт лежит в Jupyter-ноутбуке у data scientist’а, который уже уволился. Данные для ретрейна собирают руками из трёх разрозненных систем. Примерно половина проектов приходят к нам с этой болью. Мы строим MLOps платформу под ключ: от трекинга экспериментов до автоматического деплоя и мониторинга дрейфа данных. Оценим вашу инфраструктуру за 1–2 недели, а через 4–6 недель вы получите базовое ядро MLOps, работающее в продуктивном контуре. Наша команда — 10+ лет опыта в ML-инфраструктуре, более 50 внедрений.

Experiment tracking и воспроизводимость

Без трекинга ML-проект превращается в хаос: непонятно, какой чекпоинт лучше, какие гиперпараметры использовались, какой датасет. Воспроизвести результат через месяц — квест.

MLflow — open source стандарт для трекинга. Логирует параметры, метрики, артефакты (модели, графики) и код. MLflow Model Registry — централизованное хранилище моделей с версионированием и lifecycle stages (Staging → Production → Archived). Деплой через MLflow Serving или интеграция с внешними системами.

Типичная инициализация в коде:

import mlflow

mlflow.set_experiment("fraud-detection-v2")
with mlflow.start_run():
    mlflow.log_params({"learning_rate": 3e-4, "batch_size": 64, "epochs": 10})
    mlflow.log_metric("val_f1", val_f1, step=epoch)
    mlflow.pytorch.log_model(model, "model")

Это минимум. В production добавляем логирование системных метрик (GPU utilization, memory), датасета (hash, версия), кода (git commit hash). Weights & Biases — более богатый UI, collaboration features, sweep для hyperparameter optimization. MLflow — для on-premise deployment без внешних зависимостей.

DVC (Data Version Control) — версионирование данных и моделей поверх git. Данные хранятся в S3/GCS/Azure Blob, в git — только метаданные (хэши). dvc repro воспроизводит весь пайплайн от сырых данных до метрик.

Как обеспечить воспроизводимость обучения? Фиксируйте random seeds (torch.manual_seed, numpy.random.seed, random.seed) и записывайте их в метаданные эксперимента. Без этого дебаггинг нерегулярных результатов — боль. Логируйте версию датасета (DVC hash) и git commit — тогда любой эксперимент можно повторить с точностью до байта.

Оркестрация пайплайнов: Kubeflow, Airflow, Prefect

Когда нужен оркестратор пайплайнов? Скрипт обучения на 100 строк в cron — нормально для простых задач. Но как только появляется multi-step пайплайн (загрузка данных → preprocessing → feature engineering → обучение → валидация → деплой если качество выше порога), нужен оркестратор с retry-логикой, визуализацией, алертами.

Kubeflow — Kubernetes-native оркестратор для ML (см. Wikipedia). Каждый шаг — Docker-контейнер. Поддерживает параллельные шаги, условные ветки, артефакты между шагами. Интегрируется с Katib (AutoML), KServe (serving), Feast (feature store).

Apache Airflow — более общий DAG-оркестратор. Широкая экосистема операторов (S3, Spark, DBT, Kubernetes). Проще развернуть, если уже есть Airflow в компании.

Prefect / Metaflow — меньше boilerplate. Prefect 2.x с декораторами @flow и @task — быстрый старт для небольших команд.

Типичная архитектура обучающего пайплайна на Kubeflow:

  1. Data ingestion component — забирает данные из S3/БД, валидирует схему через Great Expectations
  2. Preprocessing component — трансформации, normalization, train/val/test split
  3. Training component — обучение на GPU, логирование в MLflow
  4. Evaluation component — вычисление метрик, сравнение с baseline в Model Registry
  5. Conditional deployment — деплой только если новая модель лучше текущей на >2% F1

Каждый component — отдельный Docker-образ. Пайплайн версионируется в git. Запуск по расписанию (ретрейнинг раз в неделю на новых данных) или вручную.

Model Registry и управление жизненным циклом

Model Registry — не просто хранилище чекпоинтов. Это централизованная система, которая знает:

  • Какая модель сейчас в продакшене (и с какими метриками)
  • История всех версий с параметрами обучения
  • Метаданные: датасет, git commit, результаты валидации
  • Lifecycle stage: None → Staging → Production → Archived

MLflow Model Registry — стандарт. Для enterprise — Vertex AI Model Registry (GCP), SageMaker Model Registry (AWS), Azure ML Model Registry.

Продвижение модели через стейджи: автоматически переводим модель в Staging после успешного прохождения eval, затем ручное или автоматическое (при A/B тесте) продвижение в Production. Rollback — переключение на предыдущую Production-версию за секунды.

Serving: от FastAPI до Triton Inference Server

Простой случай. FastAPI + PyTorch/ONNX на одном сервере — 80% production ML deployments именно так. Достаточно для большинства задач с нагрузкой до 100 req/s.

from fastapi import FastAPI
import onnxruntime as ort

app = FastAPI()
session = ort.InferenceSession("model.onnx", providers=["CUDAExecutionProvider"])

@app.post("/predict")
async def predict(request: PredictRequest):
    inputs = preprocess(request.text)
    outputs = session.run(None, {"input_ids": inputs})
    return {"label": postprocess(outputs)}

Triton Inference Server — production-стандарт для высоких нагрузок (500+ req/s). Dynamic batching, concurrent model execution, model ensemble. Поддерживает TensorRT, ONNX, PyTorch TorchScript, TensorFlow SavedModel.

KServe — Kubernetes-native ML serving с autoscaling, canary deployments, A/B testing из коробки. Scale-to-zero для неактивных моделей — экономия на инфраструктуре до 40% (более 1.2 млн рублей в год для проекта с 10 моделями).

Мониторинг: data drift, model drift, инфраструктурные метрики

Мониторинг — то, что обычно делают в последнюю очередь и о чём жалеют в первую. Три уровня.

Инфраструктурный мониторинг. Latency (P50/P95/P99), throughput (req/s), error rate (4xx, 5xx), GPU/CPU utilization. Prometheus + Grafana — стандарт. Алерт при P99 latency > threshold или error rate > 1%.

Data drift мониторинг. Распределение входных данных меняется со временем. Детектируем через PSI (Population Stability Index) для числовых признаков: PSI > 0.2 — сильный дрейф. Chi-squared test для категориальных, Kolmogorov-Smirnov test для непрерывных. Evidently AI — open source библиотека с готовыми дрейф-тестами.

Model drift мониторинг. Если есть ground truth с задержкой (например, через неделю знаем конверсию) — мониторим реальные метрики. Если нет — surrogate метрики: распределение prediction scores, доля confident predictions.

Alerting. Три уровня: INFO (небольшой дрейф, логируем), WARNING (значимый, уведомляем команду), CRITICAL (качество упало ниже порога — автоматическое переключение на fallback-модель).

Почему важен мониторинг дрейфа данных? Без него вы узнаёте о деградации модели только по жалобам пользователей или звенящему SLA. Алерт о дрейфе позволяет переобучить модель заранее, до того как ошибки начнут приносить убытки. В одном из наших проектов мониторинг PSI выявил дрейф через 2 дня после изменения источника данных — это спасло кампанию с бюджетами на 2 млн рублей.

Типичная ошибка Последствия Решение
Отсутствие версионирования данных Невоспроизводимость экспериментов Внедрить DVC или аналоги
Ручной деплой моделей Ошибки человеческого фактора, долгий rollback Автоматизировать CI/CD пайплайн
Мониторинг только по бизнес-метрикам Позднее обнаружение дрейфа Добавить data drift мониторинг (PSI, KS)

Feature Store

Feature Store решает проблему training-serving skew. Если preprocessing во время обучения и инференса реализован в двух разных местах — расхождение неизбежно.

Когда нужен Feature Store?

  • Несколько моделей используют одни и те же признаки
  • Признаки вычисляются из потоковых данных (real-time)
  • Большая команда с разными людьми на feature engineering и model training

Feast — open source Feature Store. Офлайн store (S3 + Parquet) для обучения, онлайн store (Redis, DynamoDB) для low-latency инференса. Feature definitions как код, materialization job синхронизирует офлайн → онлайн.

Tecton (коммерческий), Vertex AI Feature Store (GCP), SageMaker Feature Store (AWS) — managed варианты с меньшим ops overhead.

CI/CD для ML

ML CI/CD — обычный CI/CD плюс специфичные ML-шаги.

ML-специфичные checks в CI:

  • Проверка воспроизводимости: запустить обучение с фиксированным seed, результат должен совпадать
  • Data validation: Great Expectations или Pandera на schema/distribution checks
  • Model performance check: автоматический eval на holdout, блокировать merge если деградация > порога
  • Latency regression test: inference должен укладываться в SLA

GitOps для деплоя. Merge в main → CI запускает обучение → eval → если проходит → автоматический деплой в Staging → smoke tests → ручное продвижение в Production или автоматическое при успешном canary.

Инструменты: GitHub Actions / GitLab CI для CI, ArgoCD для GitOps-деплоя на Kubernetes.

Что входит в разработку MLOps-платформы

Мы предоставляем полный цикл работ, документацию и обучение команды.

Этап Длительность Результат
Аудит текущей инфраструктуры и data pipeline 1–2 недели Roadmap с рисками и приоритетами
Развёртывание ядра: MLflow, оркестратор, serving 4–6 недель Работающий пайплайн обучения и деплоя
Feature Store и CI/CD для ML 2–3 месяца Feature Store, автоматические retrain и деплой
Мониторинг дрейфа и алертинг 3–4 недели Дашборды, алерты, playbook по инцидентам
Обучение команды и документация 1–2 недели Runbook, политики, обучение для data scientists

Итоговый срок от аудита до полноценной MLOps-платформы: 3–5 месяцев. Также возможен поэтапный запуск: базовый уровень (трекинг + serving) за 4–6 недель.

Стоимость рассчитывается индивидуально под объём данных, количество моделей и требования к инфраструктуре. Закажите аудит MLOps-инфраструктуры — получите roadmap за 1–2 недели. Свяжитесь с нами для оценки вашего проекта — мы пришлём предварительный расчёт за 2 рабочих дня.

Обратите внимание: гарантия на архитектурные решения — 12 месяцев. Предоставляем сертификаты интеграции с основными облачными провайдерами (AWS, GCP, Azure). За время работы мы не потеряли ни одного клиента после первого внедрения — опыт 50+ успешных MLOps-проектов говорит сам за себя. Получите консультацию по построению MLOps платформы уже сегодня.