Что такое DAG в контексте Apache Airflow?

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

Ответ

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

  • Directed (Направленный): Зависимости между задачами имеют направление (от одной задачи к другой).
  • Acyclic (Ациклический): Граф не должен содержать циклов. Задача не может зависеть от самой себя ни прямо, ни через другие задачи, что гарантирует конечность выполнения.

На практике DAG в Airflow — это Python-скрипт, который:

  1. Определяет сам граф (объект DAG) с расписанием (schedule_interval) и другими параметрами.
  2. Определяет задачи (операторы, Operators), такие как BashOperator, PythonOperator, PostgresOperator.
  3. Определяет зависимости между задачами с помощью операторов >> (set_downstream) и << (set_upstream).

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

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

def process_data():
    print("Обрабатываю данные...")

with DAG(
    dag_id='my_etl_dag',
    start_date=datetime(2023, 1, 1),
    schedule_interval='@daily'
) as dag:

    extract = BashOperator(
        task_id='extract_data',
        bash_command='echo "Extracting..."'
    )

    transform = PythonOperator(
        task_id='transform_data',
        python_callable=process_data
    )

    load = BashOperator(
        task_id='load_data',
        bash_command='echo "Loading..."'
    )

    # Определение зависимостей: extract -> transform -> load
    extract >> transform >> load

Почему DAG — идеальная модель для оркестрации:

  • Визуализация: Airflow UI отображает граф задач, что позволяет легко отслеживать прогресс и зависимости.
  • Управление зависимостями: Позволяет четко описывать, какие задачи могут выполняться параллельно, а какие — только после успешного завершения других.
  • Перезапуск с контрольной точки: При сбое можно перезапустить только неудавшиеся задачи и их зависимых потомков, а не весь пайплайн с начала.