Что такое оператор в Apache Airflow?

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

Ответ

В Apache Airflow оператор — это шаблон задачи, который описывает единичную работу, которую нужно выполнить. Каждый оператор в DAG представляет собой один конкретный тип действия (например, запуск Python-скрипта, выполнение SQL-запроса или отправку оповещения). Операторы определяют что делать, в то время как как и где выполнять эту задачу определяет исполнитель (Executor).

Основные характеристики:

  • Атомарность: Оператор должен выполнять одну логическую операцию.
  • Идемпотентность: Повторный запуск оператора с теми же параметрами должен давать тот же результат и не вызывать побочных эффектов. Это ключевое свойство для надежности.

Популярные встроенные операторы:

  • BashOperator: Выполняет bash-команду или скрипт.
  • PythonOperator: Вызывает произвольную Python-функцию.
  • EmailOperator: Отправляет email.
  • SimpleHttpOperator: Выполняет HTTP-запрос.
  • PostgresOperator, MySqlOperator: Выполняют SQL-запрос в соответствующей БД.
  • DockerOperator: Запускает команду внутри Docker-контейнера.
  • KubernetesPodOperator: Запускает pod в Kubernetes.

Пример DAG с операторами:

from airflow import DAG
from airflow.operators.bash import BashOperator
from airflow.operators.python import PythonOperator
from datetime import datetime

def process_data():
    # Логика обработки данных
    print("Обработка данных завершена")

# Определение DAG
with DAG(
    dag_id='example_etl_dag',
    start_date=datetime(2023, 1, 1),
    schedule_interval='@daily'
) as dag:

    # Оператор 1: Загрузка данных
    download_task = BashOperator(
        task_id='download_dataset',
        bash_command='curl -o /tmp/data.csv https://example.com/data.csv'
    )

    # Оператор 2: Обработка данных
    process_task = PythonOperator(
        task_id='process_data',
        python_callable=process_data
    )

    # Оператор 3: Отправка уведомления
    notify_task = BashOperator(
        task_id='send_notification',
        bash_command='echo "ETL пайплайн успешно выполнен" | mail -s "Airflow Alert" admin@example.com'
    )

    # Определение зависимостей задач
    download_task >> process_task >> notify_task

Для обмена небольшими объемами данных между операторами в рамках одного DAG используется механизм XCom. Для более сложных сценариев интеграции данные следует хранить во внешних системах (базы данных, облачное хранилище, очереди сообщений).