Ответ
DAG (Directed Acyclic Graph, направленный ациклический граф) — это фундаментальная концепция и структурная единица в Apache Airflow. Он определяет набор задач и их зависимостей, которые необходимо выполнить.
- Directed (Направленный): Зависимости между задачами имеют направление (от одной задачи к другой).
- Acyclic (Ациклический): Граф не должен содержать циклов. Задача не может зависеть от самой себя ни прямо, ни через другие задачи, что гарантирует конечность выполнения.
На практике DAG в Airflow — это Python-скрипт, который:
- Определяет сам граф (объект
DAG) с расписанием (schedule_interval) и другими параметрами. - Определяет задачи (операторы,
Operators), такие какBashOperator,PythonOperator,PostgresOperator. - Определяет зависимости между задачами с помощью операторов
>>(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 отображает граф задач, что позволяет легко отслеживать прогресс и зависимости.
- Управление зависимостями: Позволяет четко описывать, какие задачи могут выполняться параллельно, а какие — только после успешного завершения других.
- Перезапуск с контрольной точки: При сбое можно перезапустить только неудавшиеся задачи и их зависимых потомков, а не весь пайплайн с начала.