Насколько детально вы работали с Apache AirFlow?

«Насколько детально вы работали с Apache AirFlow?» — вопрос из категории MLOps и деплой моделей, который задают на 26% собеседований Data Scientist / ML Инженер. Ниже — развёрнутый ответ с разбором ключевых моментов.

Ответ

Я использовал Apache AirFlow для оркестрации и мониторинга ML-пайплайнов и ETL-процессов в продакшене. Мой опыт включает:

1. Разработка DAG (Directed Acyclic Graph):

  • Создание тасков с использованием различных операторов: PythonOperator для ML-скриптов, BashOperator для shell-команд, DockerOperator для запуска в контейнерах.
  • Определение зависимостей между тасками с помощью битовых операторов >> и <<.
  • Использование XCom для передачи небольших данных (например, путей к файлам, ID запуска) между тасками.

2. Пример DAG для ежедневного переобучения модели:

from airflow import DAG
from airflow.operators.python import PythonOperator
from datetime import datetime, timedelta

def train_model(**context):
    # Логика обучения модели
    model_id = "model_2023_10_01"
    context['ti'].xcom_push(key='model_id', value=model_id)

def validate_model(**context):
    pulled_model_id = context['ti'].xcom_pull(task_ids='train', key='model_id')
    print(f"Validating model: {pulled_model_id}")

default_args = {
    'owner': 'data_team',
    'retries': 2,
    'retry_delay': timedelta(minutes=5),
}

with DAG(
    dag_id='daily_model_retraining',
    default_args=default_args,
    start_date=datetime(2023, 10, 1),
    schedule_interval='0 2 * * *', # Запуск каждый день в 02:00
    catchup=False
) as dag:

    train_task = PythonOperator(
        task_id='train',
        python_callable=train_model,
        provide_context=True
    )

    validate_task = PythonOperator(
        task_id='validate',
        python_callable=validate_model,
        provide_context=True
    )

    train_task >> validate_task # Определение порядка выполнения

3. Администрирование и эксплуатация:

  • Настройка расписаний (schedule_interval) и триггеров.
  • Конфигурация повторных попыток (retries), таймаутов и SLA.
  • Мониторинг выполнения через Web UI, анализ логов, настройка алертов в Slack/Telegram при сбоях.
  • Работа с пулами (pools) для ограничения параллельного выполнения ресурсоемких задач.
  • Развертывание AirFlow с использованием CeleryExecutor для распределенного выполнения задач на нескольких воркерах.