Представьте: у вас есть ETL-процесс, который каждое утро извлекает данные из трех разных источников, очищает их, преобразует и загружает в базу данных. Вчера вы запускали его вручную, а сегодня забыли. Данные не обновились, аналитики недовольны, а вы тратите часы на ручное восстановление.
Именно для таких случаев существует Apache Airflow и его главный инструмент — DAG. Это не просто способ автоматизировать задачи, а целая философия оркестрации данных, которая избавит вас от рутины и ошибок.
В этой статье разберем, что такое DAG в Airflow и как создавать свои графы задач, настраивать расписание и избегать типичных ошибок — от неправильного catchup до нарушения идемпотентности.
Что такое DAG и зачем он нужен в Airflow
DAG (Directed Acyclic Graph) — направленный ациклический граф. Звучит сложно, но на практике это просто схема вашего рабочего процесса:
- Directed (направленный) — задачи выполняются в определенном порядке;
- Acyclic (ациклический) — нет циклов, задача не может зависеть сама от себя;
- Graph (граф) — структура из узлов (задач) и ребер (зависимостей).
То есть DAG — это граф, который говорит Airflow, какие задачи, в каком порядке и когда выполнять.
Основные понятия: задачи, операторы и зависимости
В экосистеме Airflow есть три ключевых элемента:
Задачи — это атомарные единицы работы. Каждая задача выполняет конкретное действие: запуск скрипта, выполнение SQL-запроса, вызов API.
Операторы — это предопределенные шаблоны для создания задач. Airflow предоставляет множество встроенных операторов:
- BashOperator — выполнение команд bash;
- PythonOperator — выполнение Python-функций;
- SQLOperator — выполнение SQL-запросов;
- HttpOperator — выполнение HTTP-запросов;
- DockerOperator — запуск контейнеров Docker.
Зачем нужны DAG в оркестрации данных
DAG решают несколько важных задач в работе с данными:
- автоматизация рутинных процессов — ETL-пайплайны, очистка данных, генерация отчетов;
- управление зависимостями — задачи выполняются в правильном порядке;
- мониторинг и логирование — встроенный веб-интерфейс для отслеживания статуса;
- обработка ошибок;
- масштабируемость — распределенное выполнение задач на разных воркерах.
DAG не выполняет задачи сам по себе. Он только описывает, что и как нужно делать. Выполнением занимаются воркеры Airflow.
Структура Airflow DAG: из чего состоит граф задач
Каждый Airflow DAG имеет определенную структуру и набор параметров, которые определяют его поведение. Давайте разберем их по порядку.
Параметры DAG
При создании DAG необходимо задать ключевые параметры:
from airflow import DAGfrom datetime import datetime, timedelta
default_args = { 'owner': 'data_team', 'depends_on_past': False, 'start_date': datetime(2026, 6, 15), # дата начала выполнения DAG 'email': ['example@example.com'], 'email_on_failure': True, 'email_on_retry': False, 'retries': 3, 'retry_delay': timedelta(minutes=5),}
dag = DAG( 'etl_pipeline', default_args=default_args, schedule='@daily', catchup=False, tags=['etl', 'production'])
Ключевые параметры:
Параметр depends_on_past определяет, зависит ли задача от успешного выполнения предыдущего запуска этой же задачи. По умолчанию False.
Операторы в Airflow: Python и Bash
Операторы определяют тип задачи. Рассмотрим два самых популярных.
BashOperator — для выполнения команд оболочки:
from airflow.operators.bash import BashOperator bash_task = BashOperator( task_id='bash_script', bash_command='echo "Hello, Airflow!" && date', dag=dag)
PythonOperator — для выполнения Python-функций. Это основной способ работы с Apache Airflow:
from airflow.operators.python import PythonOperator
def process_data(**context): print("Processing data...") # Логика обработки данных return {'status': 'success'}
python_task = PythonOperator( task_id='process_data', python_callable=process_data, dag=dag)
Зависимости между задачами
Задачи в DAG связаны зависимостями, которые определяют порядок выполнения. Airflow использует операторы побитового сдвига для задания зависимостей:
# Последовательное выполнениеtask1 >> task2 >> task3 # Параллельное выполнениеtask1 >> [task2, task3] >> task4 # Сложные зависимости[task1, task2] >> task3 >> [task4, task5]
По умолчанию задача начнет выполнение только если все предыдущие задачи завершились успешно. Это поведение можно изменить через параметр trigger_rule.
Создаем свой первый Airflow DAG: пошаговое руководство
Рассмотрим простой пример: у нас есть таблица со студентами, и нам нужно каждый день отбирать отличников (тех, у кого средний балл выше 4,5) и сохранять их в отдельную таблицу.
Для работы нам понадобится SQLite — эта база данных встроена в Python, так что ничего устанавливать не нужно. Импортируем необходимые модули:
import sqlite3from datetime import datetime, timedeltafrom airflow import DAGfrom airflow.operators.python import PythonOperator
Задаем базовые параметры для всех задач:
default_args = { 'owner': 'data_engineer', 'depends_on_past': False, 'start_date': datetime(2026, 1, 1), 'email_on_failure': False, 'email_on_retry': False, 'retries': 2, 'retry_delay': timedelta(minutes=5),}
dag = DAG( 'student_filter_pipeline', default_args=default_args, description='Фильтрация студентов с высоким средним баллом', schedule='@daily', catchup=False, tags=['students', 'tutorial'])
Перед запуском DAG создадим таблицу со студентами. Это можно сделать один раз вручную:
# Создаем тестовую базу данныхconn = sqlite3.connect('/opt/airflow/data/students.db')cursor = conn.cursor()
# Создаем таблицу студентовcursor.execute(''' CREATE TABLE IF NOT EXISTS students ( id INTEGER PRIMARY KEY, name TEXT, average_score REAL )''')
# Добавляем тестовые данныеstudents_data = [ ('Иванов Иван', 4.8), ('Петров Петр', 3.5), ('Сидорова Анна', 4.9), ('Козлов Дмитрий', 4.2), ('Смирнов Максим', 4.6), ('Яковцев Илья', 3.9),]
cursor.executemany('INSERT INTO students (name, average_score) VALUES (?, ?)', students_data)conn.commit()conn.close()
Создаем три задачи: Extract (извлечение), Filter (фильтрация), Load (загрузка).
Задача 1: Extract — извлекаем данные из таблицы студентов
def extract_students(**context): """Извлекает данные из таблицы студентов""" conn = sqlite3.connect('/opt/airflow/data/students.db') cursor = conn.cursor() # Получаем всех студентов cursor.execute('SELECT id, name, average_score FROM students') students = cursor.fetchall() conn.close() # Передаем данные в следующую задачу через XCom context['ti'].xcom_push(key='students_data', value=students) print(f"Извлечено {len(students)} студентов")
extract_task = PythonOperator( task_id='extract_students', python_callable=extract_students, dag=dag)
Задача 2: Filter — оставляем только отличников
def filter_top_students(**context): """Фильтрует студентов: оставляем только тех, у кого средний балл > 4.5""" # Получаем данные из предыдущей задачи ti = context['ti'] students_data = ti.xcom_pull(task_ids='extract_students', key='students_data') # Фильтрация: средний балл выше 4.5 top_students = [ student for student in students_data if student[2] > 4.5 # student[2] - это average_score ] # Передаем отфильтрованные данные ti.xcom_push(key='top_students_data', value=top_students) print(f"Отфильтровано {len(top_students)} отличников")
filter_task = PythonOperator( task_id='filter_top_students', python_callable=filter_top_students, dag=dag)
Задача 3: Load — записываем отличников в отдельную таблицу
def load_top_students(**context): """Загружает отличников в таблицу top_students""" ti = context['ti'] top_students_data = ti.xcom_pull(task_ids='filter_top_students', key='top_students_data') conn = sqlite3.connect('/opt/airflow/data/students.db') cursor = conn.cursor() # Создаем таблицу для отличников (если ее нет) cursor.execute(''' CREATE TABLE IF NOT EXISTS top_students ( id INTEGER PRIMARY KEY, name TEXT, average_score REAL ) ''') # Очищаем таблицу перед записью (идемпотентность!) cursor.execute('DELETE FROM top_students') # Записываем новых отличников cursor.executemany( 'INSERT INTO top_students (id, name, average_score) VALUES (?, ?, ?)', top_students_data ) conn.commit() conn.close() print(f"Загружено {len(top_students_data)} отличников в таблицу top_students")
load_task = PythonOperator( task_id='load_top_students', python_callable=load_top_students, dag=dag)
Как это работает на практике
Давайте посмотрим, что происходит при запуске этого DAG:
- Extract — подключается к базе данных и извлекает всех студентов из таблицы students.
- Filter — получает список студентов и оставляет только тех, у кого average_score > 4.5.
- Load — создает, если нужно, таблицу top_students, очищает ее и записывает отфильтрованных студентов.
В результате в таблице top_students окажутся только отличники:
- Иванов Иван (4.8);
- Сидорова Анна (4.9);
- Смирнов Максим (4.6).
Запуск и мониторинг через веб-интерфейс
После создания DAG-файла:
- Разместите файл в директории dags/ (обычно ~/airflow/dags/ или /opt/airflow/dags/).
- Обновите список DAG (если нужно):
airflow dags list
- Откройте веб-интерфейс (по умолчанию http://localhost:8080).
- Найдите ваш DAG в списке и активируйте его переключателем.
- Запустите вручную (для тестирования) или дождитесь расписания.
Мониторинг выполнения:
- Graph View — визуализация графа задач;
- Tree View — древовидное представление запусков;
- Logs — логи выполнения каждой задачи;
- Gantt Chart — временная диаграмма выполнения.
Концепции планирования в Airflow: schedule и catchup
Catchup в Airflow — механизм, при котором Airflow наверстывает пропущенные запуски DAG.
Если вы создали DAG с start_date в прошлом и catchup=True, Airflow создаст и выполнит запуски за весь период от start_date до текущего момента.
Сценарий:
- start_date: 2026-06-01;
- schedule: @daily;
- Текущая дата: 2026-06-15.
При catchup=True Airflow создаст и выполнит 14 запусков DAG (за каждый день с 1 по 14 июня).
При catchup=False Airflow создаст только один запуск для текущего интервала.
Настройка schedule: cron, timedelta и timetable
Airflow поддерживает несколько способов задания расписания через параметр schedule:
1. Cron-выражения:
# Каждый день в 03:00schedule='0 3 * * *' # Каждый понедельник в 09:00schedule='0 9 * * 1' # Каждые 15 минутschedule='*/15 * * * *'
2. Предустановки (presets):
3. Timedelta:
from datetime import timedelta # Каждые 30 минутschedule=timedelta(minutes=30) # Каждые 6 часовschedule=timedelta(hours=6)
4. None (ручной запуск):
# DAG не запускается по расписанию, только вручнуюschedule=None
Task instance в Airflow: мониторинг выполнения задач
Task instance airflow — это экземпляр задачи для конкретного запуска DAG. Каждый DAG run создает task instances для всех задач в DAG.
Каждая задача может находиться в одном из следующих состояний:
Удобнее всего отслеживать статус задач в веб-интерфейсе Airflow. С помощью него можно:
- просмотреть статус всех task instances для конкретного DAG run;
- проверить логи выполнения каждой task instance;
- перезапустить task instances;
- отметить task instances как success или failed вручную.
Распространенные ошибки и как их избежать
Airflow DAGs: коротко о главном
- DAG — направленный ациклический граф, определяющий структуру и расписание задач в Apache Airflow.
- Операторы — это шаблоны для создания задач (BashOperator, PythonOperator и др.).
- Зависимости определяют порядок выполнения задач через операторы >>.
- Schedule — поддерживаются cron-выражения, timedelta и кастомные Timetables.
- Catchup и Backfill — механизмы для выполнения пропущенных запусков.
- Task instance — экземпляр задачи для конкретного dag run.
- Идемпотентность критически важна для надежности пайплайнов.
- Мониторинг работает через встроенный веб-интерфейс для отслеживания статуса задач.
Теперь вы не просто знаете, что такое DAG в Airflow, а понимаете, как создавать надежные, масштабируемые и поддерживаемые пайплайны — как настоящий data-инженер.
