Был ли опыт работы с оркестратором?

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

Ответ

Да, я работал с Apache Airflow для оркестрации ETL-пайплайнов и задач по обработке данных. Основная задача — создание, планирование и мониторинг сложных рабочих процессов, состоящих из множества зависимых задач.

Мой опыт с Airflow включает:

  • Разработка DAG (Directed Acyclic Graph): Описание зависимостей между задачами (тасками) с помощью Python.
  • Создание кастомных операторов: Когда встроенных (PythonOperator, BashOperator) было недостаточно.
  • Настройка расписания: Использование cron-выражений или временных интервалов для запуска пайплайнов.
  • Обработка ошибок и повторные попытки: Конфигурация retries, retry_delay, email_on_failure.
  • Передача данных между задачами: Использование XCom для обмена небольшими сообщениями.
  • Работа с переменными и подключениями: Хранение конфигурации в Airflow Variables и секретов (паролей, ключей) в Connections.
  • Мониторинг: Использование веб-интерфейса Airflow для отслеживания статусов DAG и задач, просмотра логов.

Пример простого DAG для ETL-пайплайна:

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

def extract(**context):
    # Логика извлечения данных (например, из API или БД)
    data = [1, 2, 3, 4, 5]
    # Отправляем данные в XCom для следующей задачи
    context['ti'].xcom_push(key='raw_data', value=data)
    print(f"Extracted data: {data}")

def transform(**context):
    # Получаем данные из предыдущей задачи через XCom
    pulled_data = context['ti'].xcom_pull(key='raw_data', task_ids='extract_task')
    # Простая трансформация
    transformed_data = [x * 2 for x in pulled_data]
    context['ti'].xcom_push(key='transformed_data', value=transformed_data)
    print(f"Transformed data: {transformed_data}")

def load(**context):
    transformed_data = context['ti'].xcom_pull(key='transformed_data', task_ids='transform_task')
    # Логика загрузки (например, в другую БД или файл)
    print(f"Loading data: {transformed_data} to destination...")

# Определение аргументов DAG по умолчанию
default_args = {
    'owner': 'data_team',
    'depends_on_past': False,
    'retries': 2,
    'retry_delay': timedelta(minutes=5),
    'email_on_failure': True,
}

# Создание DAG
with DAG(
    'simple_etl_pipeline',
    default_args=default_args,
    description='A simple ETL pipeline',
    schedule_interval='@daily',  # Запуск раз в день
    start_date=datetime(2023, 1, 1),
    catchup=False,  # Не запускать пропущенные интервалы
    tags=['etl', 'example'],
) as dag:

    start = DummyOperator(task_id='start')

    extract_task = PythonOperator(
        task_id='extract_task',
        python_callable=extract,
        provide_context=True,
    )

    transform_task = PythonOperator(
        task_id='transform_task',
        python_callable=transform,
        provide_context=True,
    )

    load_task = PythonOperator(
        task_id='load_task',
        python_callable=load,
        provide_context=True,
    )

    end = DummyOperator(task_id='end')

    # Определение порядка выполнения задач
    start >> extract_task >> transform_task >> load_task >> end

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