Баннер мобильный (3) Пройти тест

Airflow DAGs: как создавать и запускать графы задач в Apache Airflow

Разбираем основы оркестрации данных на примерах

Инструкция

8 июля 2026

Поделиться

Скопировано
Airflow DAGs: как создавать и запускать графы задач в Apache Airflow

Содержание

    Представьте: у вас есть ETL-процесс, который каждое утро извлекает данные из трех разных источников, очищает их, преобразует и загружает в базу данных. Вчера вы запускали его вручную, а сегодня забыли. Данные не обновились, аналитики недовольны, а вы тратите часы на ручное восстановление.

    Именно для таких случаев существует Apache Airflow и его главный инструмент — DAG. Это не просто способ автоматизировать задачи, а целая философия оркестрации данных, которая избавит вас от рутины и ошибок.

    В этой статье разберем, что такое DAG в Airflow и как создавать свои графы задач, настраивать расписание и избегать типичных ошибок — от неправильного catchup до нарушения идемпотентности.

    Что такое DAG и зачем он нужен в Airflow

    DAG (Directed Acyclic Graph) — направленный ациклический граф. Звучит сложно, но на практике это просто схема вашего рабочего процесса:

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

    То есть DAG — это граф, который говорит Airflow, какие задачи, в каком порядке и когда выполнять.

    Основные понятия: задачи, операторы и зависимости

    В экосистеме Airflow есть три ключевых элемента:

    Понятие
    Что это
    Пример
    DAG
    Контейнер для всех задач, который определяет структуру и расписание
    dag = DAG(‘my_dag’, …)
    Task (задача)
    Единица работы внутри DAG
    Извлечение данных из API
    Operator
    Шаблон для создания задач определенного типа
    BashOperator, PythonOperator

    Задачи — это атомарные единицы работы. Каждая задача выполняет конкретное действие: запуск скрипта, выполнение 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'])

    Ключевые параметры:

    Параметр
    Что делает
    Пример значения
    dag_id
    Уникальный идентификатор DAG
    ‘etl_pipeline’
    start_date
    Дата начала выполнения DAG
    datetime(2026, 6, 15)
    schedule
    Интервал запуска
    ‘@daily’, timedelta(hours=1)
    catchup
    Выполнять пропущенные запуски
    True / False
    default_args
    Аргументы по умолчанию для всех задач
    Словарь с настройками
    tags
    Теги для группировки DAG
    [‘etl’, ‘production’]

    Параметр depends_on_past определяет, зависит ли задача от успешного выполнения предыдущего запуска этой же задачи. По умолчанию False.

    Примечание: В версиях Airflow до 2.4 использовался параметр schedule_interval. Начиная с Airflow 2.4, он заменен на schedule. Если вы работаете с более старой версией, используйте schedule_interval — он все еще поддерживается для обратной совместимости.

    Операторы в 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:

    1. Extract — подключается к базе данных и извлекает всех студентов из таблицы students.
    2. Filter — получает список студентов и оставляет только тех, у кого average_score > 4.5.
    3. Load — создает, если нужно, таблицу top_students, очищает ее и записывает отфильтрованных студентов.

    В результате в таблице top_students окажутся только отличники:

    • Иванов Иван (4.8);
    • Сидорова Анна (4.9);
    • Смирнов Максим (4.6).

    Запуск и мониторинг через веб-интерфейс

    После создания DAG-файла:

    1. Разместите файл в директории dags/ (обычно ~/airflow/dags/ или /opt/airflow/dags/).
    2. Обновите список DAG (если нужно):
    airflow dags list
    1. Откройте веб-интерфейс (по умолчанию http://localhost:8080).
    2. Найдите ваш DAG в списке и активируйте его переключателем.
    3. Запустите вручную (для тестирования) или дождитесь расписания.

    Мониторинг выполнения:

    • 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):

    Предустановка
    Значение
    Описание
    @once
    Единожды
    Запуск только один раз
    @hourly
    0 * * * *
    Каждый час
    @daily
    0 0 * * *
    Каждый день в полночь
    @weekly
    0 0 * * 0
    Каждое воскресенье
    @monthly
    0 0 1 * *
    Первый день месяца
    @yearly
    0 0 1 1 *
    Первый день года

    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.

    Каждая задача может находиться в одном из следующих состояний:

    Статус
    Описание
    None
    Задача еще не запущена
    scheduled
    Задача запланирована к выполнению
    queued
    Задача в очереди на выполнение
    running
    Задача выполняется
    success
    Задача успешно завершена
    failed
    Задача завершилась с ошибкой
    upstream_failed
    Предыдущая задача завершилась с ошибкой

    Удобнее всего отслеживать статус задач в веб-интерфейсе Airflow. С помощью него можно:

    • просмотреть статус всех task instances для конкретного DAG run;
    • проверить логи выполнения каждой task instance;
    • перезапустить task instances;
    • отметить task instances как success или failed вручную.

    Распространенные ошибки и как их избежать

    Ошибка
    Причина
    Решение
    DAG не запускается
    Неправильный start_date
    Убедитесь, что start_date — это дата в прошлом
    Слишком много запусков
    catchup=True при слишком старом start_date
    Используйте catchup=False
    Задачи выполняются в неправильном порядке
    Не заданы зависимости
    Проверьте операторы >>
    Дубликаты данных
    Неидемпотентные задачи
    Перезаписывайте данные, а не добавляйте
    Задача падает без понятной ошибки
    Нет логирования
    Добавьте print(), а лучше используйте logging

    Airflow DAGs: коротко о главном

    • DAG — направленный ациклический граф, определяющий структуру и расписание задач в Apache Airflow.
    • Операторы — это шаблоны для создания задач (BashOperator, PythonOperator и др.).
    • Зависимости определяют порядок выполнения задач через операторы >>.
    • Schedule — поддерживаются cron-выражения, timedelta и кастомные Timetables.
    • Catchup и Backfill — механизмы для выполнения пропущенных запусков.
    • Task instance — экземпляр задачи для конкретного dag run.
    • Идемпотентность критически важна для надежности пайплайнов.
    • Мониторинг работает через встроенный веб-интерфейс для отслеживания статуса задач.

    Теперь вы не просто знаете, что такое DAG в Airflow, а понимаете, как создавать надежные, масштабируемые и поддерживаемые пайплайны — как настоящий data-инженер.

    Инструкция

    Поделиться

    Скопировано
    0 комментариев
    Комментарии