Скидки до 55% и курс по ИИ в подарок 2 дня 09 :56 :09 Выбрать курс
Код
#статьи

Что такое Apache Airflow и как с ней работать

Автоматизируем рабочие процессы для обработки данных.

Иллюстрация: Polina Vari для Skillbox Media

Apache Airflow — это платформа с открытым исходным кодом для автоматизации процессов обработки данных. С её помощью можно запускать связанные задачи в нужном порядке, задавать расписание и контролировать их выполнение.

В этой статье мы простыми словами разберём, как устроена Apache Airflow, из каких компонентов она состоит, где её применяют и как создать и запустить первый пайплайн.

Содержание


Кому и для чего нужна Apache Airflow

Apache Airflow нужна там, где работу с данными приходится выполнять регулярно и по определённому сценарию. Например, для автоматизации ETL- и ELT-процессов. В ETL данные сначала извлекают, затем преобразуют и загружают в хранилище. В ELT их сначала загружают, а преобразуют уже внутри хранилища.

Такие последовательности действий называют пайплайнами данных. Airflow позволяет описать пайплайн в коде, задать порядок выполнения задач и автоматически запускать их по расписанию или при определённых условиях. Такое управление пайплайнами упрощает работу со сложными процессами, состоящими из нескольких этапов. Если один из них завершится с ошибкой, Airflow покажет, где именно возникла проблема, и поможет перезапустить нужную часть процесса. За счёт этого специалистам легче осуществлять контроль ошибок.

Чаще всего с Airflow работают дата-инженеры. Они строят пайплайны, которые собирают данные из баз, API и других источников, преобразуют их и передают в аналитические системы или хранилища. Также инструмент используют разработчики, аналитики и специалисты по машинному обучению.

Например, интернет-магазин может каждую ночь собирать данные о просмотрах товаров, заказах и остатках на складах. Airflow последовательно выполнит запуск задач по расписанию: выгрузит информацию из разных систем IT-системы, очистит её, объединит и загрузит в хранилище. Утром аналитики смогут построить отчёты по продажам, а рекомендательная система — использовать свежие данные о поведении покупателей для формирования индивидуальных предложений посетителям интернет-магазина.

В машинном обучении Airflow применяют похожим образом. Допустим, модель нужно регулярно переобучать на новых данных. Airflow может запустить их загрузку и подготовку, затем обучение модели, проверку её качества и следующий этап пайплайна. При этом сама Airflow не обучает модель и не обрабатывает данные вместо специализированных инструментов — она управляет последовательностью этих операций.

Поэтому Airflow правильнее воспринимать не как систему для хранения или обработки данных, а как оркестратор. Она связывает отдельные задачи в единый процесс, следит за зависимостями между ними и запускает их в нужном порядке.

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

Как работает Apache Airflow

Apache Airflow управляет задачами, которые должны выполняться в определённом порядке. Вы описываете на Python, что нужно сделать, когда запускать каждую задачу и от каких других задач она зависит, а Airflow следит за выполнением этого сценария.

Например, в пайплайне обработки данных сначала нужно получить информацию из нескольких источников, затем обработать её и объединить, а после — загрузить результат в хранилище, чтобы её могли использовать в работе. Airflow запустит задачи в нужной последовательности и не перейдёт к следующему этапу, пока не выполнены необходимые условия.

Airflow собирает данные из разных источников, отправляет на обработку и выгружает данные, готовые для анализа
Изображение: Mermaid.ai / Skillbox Media

За работу Airflow отвечают несколько компонентов:

  • DAG описывает задачи и связи между ними;
  • планировщик определяет, какие задачи и когда нужно запустить;
  • исполнитель передаёт задачи на выполнение;
  • рабочие непосредственно выполняют задачи;
  • DAG-процессор считывает и обрабатывает файлы с DAG;
  • база метаданных хранит информацию о задачах, запусках и их статусах;
  • триггер-менеджер обслуживает задачи, которые ждут наступления определённых событий;
  • веб-сервер предоставляет интерфейс для наблюдения за пайплайнами и управления ими.

Дальше разберём каждый компонент подробнее на примере пайплайна банковского скоринга.

DAG

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

Граф состоит из узлов и связей между ними. В Airflow узлы — это задачи, а связи показывают зависимости: например, задачу B можно запустить только после успешного выполнения задачи A. Граф называют направленным, потому что у этих связей есть направление, и ацикличным, потому что цепочка задач не может замкнуться сама на себя.

В Airflow DAG описывает структуру всего пайплайна и создаётся на Python. Посмотрим на пример банковского скоринга. Для этого нам потребуется описать семь задач и настроить условия проверки данных и скоринга.

Первой в пайплайне будет задача ingest: система собирает данные о потенциальном заёмщике из необходимых источников. После этого задача analyze обрабатывает их и рассчитывает показатели, которые понадобятся для оценки заявки, — например, среднемесячный доход и долговую нагрузку.

Задача check integrity проверяет полноту и согласованность данных, а также результат скоринга. Дальше пайплайн может пойти по одному из двух сценариев:

  • если обнаружены ошибки или скоринговый балл слишком низкий, система формирует причину отказа (describe integrity) и передаёт информацию сотруднику (email error);
  • если проверка пройдена, заявка получает положительный статус (save).

В конце задача report сохраняет итоговую информацию о заявке и принятом решении.

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

DAG с узлами-задачами
Иллюстрация: Mermaid.ai / Skillbox Media

Планировщик

Планировщик (Scheduler) определяет, когда запускать DAG и какие задачи внутри него готовы к выполнению. Для этого он учитывает расписание, зависимости между задачами и их текущее состояние.

Например, банковский пайплайн можно запускать по расписанию — допустим, раз в час для обработки новых заявок. Другой вариант — настроить внешний сервис так, чтобы он инициировал запуск DAG сразу после поступления заявки от клиента.

После запуска планировщик проверяет зависимости между задачами. Например, анализ данных начнётся только после того, как система соберёт сведения о заёмщике из государственных сервисов.

Сам планировщик задачи не выполняет. Он определяет, что и когда нужно запустить, а затем передаёт готовые задачи исполнителю

Исполнитель

Исполнитель (Executor) получает от планировщика задачи, готовые к запуску, и определяет, где их выполнять. В простой конфигурации все задачи могут выполняться на одной машине, а в распределённой — на разных рабочих узлах или в контейнерах.

Например, планировщик решил, что пора получить сведения о заёмщике из реестра банкротов. Исполнитель принимает задачу и передаёт её свободному рабочему, который обращается к реестру банкротов по API.

Рабочий

Рабочий (Worker) выполняет сам код задачи: обращается к внешним сервисам и базам данных, запускает скрипты, обрабатывает информацию и сохраняет результат.

В примере с банком рабочий получает задачу собрать данные о заёмщике и делает запрос к нужному источнику. В распределённой системе рабочих может быть несколько, поэтому Airflow способна одновременно обрабатывать много независимых задач.

DAG-процессор

DAG-процессор регулярно считывает Python-файлы с DAG и извлекает из них информацию о задачах, зависимостях и настройках пайплайна. Если разработчик изменит DAG банковского скоринга, процессор обнаружит обновление и обработает новый код.

База метаданных

База метаданных (Metadata Metabase) хранит служебную информацию о работе Airflow: данные о DAG и их запусках, состояния задач, расписания, подключения и другие настройки.

К этой базе обращаются разные компоненты системы. Например, планировщик узнаёт из неё, какие задачи уже завершились, а веб-интерфейс использует эти данные, чтобы показывать историю и статусы запусков.

Триггер-менеджер

Триггер-менеджер нужен для задач, которые ждут внешнего события — например, ответа API или появления новых данных. Вместо того чтобы всё время занимать рабочий процесс, задача переходит в режим ожидания, а триггер-менеджер следит за нужным событием.

В банковском скоринге это пригодится, если один из внешних сервисов отвечает не сразу. Когда данные появятся, триггер-менеджер сообщит Airflow, что задачу можно продолжить.

Веб-сервер

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

Например, сотрудник банка может открыть пайплайн скоринга и проверить, на каком этапе остановилась обработка заявки. При необходимости из интерфейса можно вручную запустить DAG или повторить выполнение отдельных задач.

Плагины

Плагины используются для расширения возможностей Airflow и добавления собственных компонентов. С их помощью можно адаптировать систему под внутреннюю инфраструктуру компании или специфические рабочие процессы.

Например, банк может создать плагин для взаимодействия с внутренним сервисом. При этом для многих популярных баз данных, облачных платформ и API в Airflow уже существуют готовые интеграции — провайдеры.

Общая схема архитектуры Apache Airflow
Иллюстрация: Mermaid.ai / Skillbox Media

Как установить и запустить Apache Airflow

В предыдущем разделе мы разобрали, как устроена Apache Airflow. Теперь перейдём к практике: установим Airflow и запустим её локально.

Будем использовать минимальную конфигурацию без реальных внешних сервисов. Их заменим моками — функциями, которые имитируют работу API, баз данных и других систем.

Для работы понадобится установленный Python. Подойдёт версия Python 3.8 или выше, но лучше использовать Python 3.10 или более новую совместимую версию.

Разбираться во всех деталях кода необязательно: примеры можно скопировать и сохранить в файлы. Чтобы было проще понять, что происходит на каждом этапе, мы добавили в код комментарии.

Шаг 1: создайте виртуальное окружение. Для этого введите в терминале:

python3 -m venv airflow_env

Если терминал не может найти команду python3, замените на python.

Шаг 2: активируйте виртуальное окружение. Для Linux/macOS введите:

source airflow_env/bin/activate

Для Windows:

airflow_env\Scripts\activate

Если окружение активировалось, то перед строкой приглашения в терминале появится (airflow_env).

Шаг 3: установите Apache Airflow вместе с совместимыми версиями зависимостей. Если делегировать это pip, то есть риск конфликтов между пакетами.

Поэтому используем официальный файл ограничений — constraints-файл. Он фиксирует набор совместимых версий и помогает получить стабильную и воспроизводимую среду.

AIRFLOW_VERSION=3.0.6 PYTHON_VERSION=$(python -c "import sys; print('.'.join(map(str, sys.version_info[:2])))") CONSTRAINT_URL="https://raw.githubusercontent.com/apache/airflow/constraints-
AIRFLOWVERSION/constraints-
{PYTHON_VERSION}.txt" pip install "apache-airflow==
AIRFLOWVERSION"--constraint"
{CONSTRAINT_URL}"

Установка займёт некоторое время. В терминале появится много строк с информацией о загрузке и установке пакетов. Дождитесь, пока снова не появится строка приглашения командной оболочки.

Шаг 4: запустите Airflow в режиме standalone. Он подходит для локального знакомства с системой: Airflow автоматически подготовит базу метаданных, создаст пользователя и запустит основные компоненты.

Для этого введите в терминал:

airflow standalone

После запуска в терминале появятся адрес веб-интерфейса, логин и пароль. Сохраните эти данные — они понадобятся на следующем шаге.

Режим standalone предназначен для локальной разработки и обучения. В продакшене компоненты Airflow обычно настраивают и запускают отдельно, а для базы метаданных используют полноценную СУБД.

Шаг 5: проверьте, что Airflow запустился. Откройте в браузере адрес http://localhost:8080. Если всё настроено правильно, появится пустой веб-интерфейс Apache Airflow. Теперь остаётся написать первый DAG.

Веб-интерфейс Apache Airflow
Скриншот: Apache Airflow / Google Chrome / Skillbox Media

Пишем конвейер задач в Apache Airflow для обработки данных

Теперь соберём собственный конвейер в Airflow. Создадим два DAG: сначала простой, с одной задачей — а затем более сложный, который будет имитировать сбор и обработку данных интернет-магазина.

Однозадачный DAG для генерации случайных чисел

На его примере посмотрим, как Airflow запускает задачу и повторяет её при ошибке. Создайте файл first_dag.py в папке ~/airflow/dags. Дополнительно подключать его к Airflow не нужно: планировщик автоматически обнаружит новый DAG.

import random  # Модуль для генерации случайных чисел
import datetime as dt  # Модуль для работы с датами и временем

from airflow.models import DAG  # Импорт класса DAG для определения графа задач
from airflow.operators.python import PythonOperator  # Импорт оператора для выполнения Python-функций

# Установим аргументы по умолчанию для задач в DAG
default_args = {
    'owner': 'airflow', 
    'start_date': dt.datetime(2025, 9, 16),  # Дата начала (DAG не запустится раньше)
    'retries': 2,  # Количество повторных попыток при ошибке
    'retry_delay': dt.timedelta(seconds=10),  # Задержка между попытками
}

# Функция задачи: имитирует бросок кубика; поднимает ошибку, если выпадает нечётное (шанс 50%)
def random_dice():
    val = random.randint(1, 6)  # Генерация случайного числа от 1 до 6
    print(f"Бросок кубика: {val}")  # Вывод результата в логи
    if val % 2 != 0:  # Проверка на нечётность
        raise ValueError(f'Odd {val}')  # Симуляция ошибки (задача провалится)

# Определение DAG: граф с одной задачей
with DAG(
    dag_id='first_dag',  # Уникальный ID DAG (виден в UI)
    schedule='@daily',  # Расписание: ежедневно; для теста можно '@once'
    default_args=default_args,  # Применение дефолтных аргументов
    catchup=False  # Не выполнять пропущенные запуски
) as dag:
    # Создание задачи: PythonOperator для вызова функции
    dice = PythonOperator(
        task_id='random_dice',  # Уникальный ID задачи
        python_callable=random_dice,  # Функция для выполнения
        dag=dag,  # Связь с DAG (опционально, но рекомендуется)
    )

Этот DAG запускается раз в сутки и содержит одну задачу — random_dice. При каждом запуске функция генерирует случайное число от 1 до 6, как при броске игрального кубика.

Если выпадает чётное число, задача успешно завершается. Если нечётное — функция вызывает ошибку ValueError, и Airflow помечает попытку как неудачную. После этого она ждёт десять секунд и запускает задачу повторно. Всего предусмотрено две повторные попытки.

Такой пример позволяет сразу увидеть один из базовых механизмов Airflow — автоматические повторные попытки. Если задача временно завершилась с ошибкой, её не обязательно перезапускать вручную: Airflow сделает это сама по заданным правилам.

Параметр catchup=False при этом не позволяет Airflow создавать запуски за все пропущенные даты, начиная с start_date. Поэтому после добавления DAG система будет выполнять только актуальные запуски по расписанию.

Проверим DAG в работе. Откройте http://localhost:8080, найдите в списке first_dag и активируйте его переключателем слева. После этого Airflow будет запускать DAG по заданному расписанию. Чтобы не ждать следующего запуска, нажмите кнопку ручного запуска справа.

Активируем и запускаем DAG
Скриншот: Apache Airflow / Google Chrome / Skillbox Media

После запуска откройте страницу DAG и перейдите к списку его запусков и задач. Там можно посмотреть статус random_dice и открыть логи выполнения.

В логах появится результат броска кубика:

  • если выпало чётное число, задача завершится успешно;
  • если нечётное — возникнет ошибка ValueError. Airflow подождёт десять секунд и попробует выполнить задачу снова. Всего она сделает до двух повторных попыток, как указано в default_args.

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

Наш первый DAG получился максимально простым: в нём всего одна задача. Его структуру можно посмотреть в представлении Graph или на странице конкретного запуска DAG — там будет только один узел random_dice.

Структура DAGа с одним действием
Скриншот: Apache Airflow / Google Chrome / Skillbox Media

Многозадачный DAG для обработки онлайн-заказов

Напишем DAG из нескольких связанных задач по обработке заказов в онлайн-магазине. Сначала Airflow должен будет получить данные о заказе, затем проверить наличие товара и в зависимости от результата выбрать одно из двух действий: обработать заказ или отправить клиенту уведомление о том, что товара нет в наличии.

Скопируйте код в новый файл и сохраните его в папке ~/airflow/dags.

from __future__ import annotations # Импорт для поддержки новой синтаксической конструкции для аннотаций типов

import pendulum # Библиотека для работы с датами и временем
import datetime as dt # Также импортируем стандартную библиотеку для работы с датами

from airflow.models.dag import DAG # Импортируем основной класс DAG
from airflow.operators.python import PythonOperator, BranchPythonOperator # Импортируем операторы, которые будут выполнять Python-функции и управлять ветвлением

# Создаём DAG, используя менеджер контекста (`with DAG(...) as dag:`)
# Это рекомендуемый способ, который автоматически регистрирует задачи в DAG.
with DAG(
    dag_id="order_processing_dag", # Уникальный идентификатор DAG
    # start_date задаёт, с какой даты Airflow начнёт планировать задачи.
    # Используем pendulum для работы с часовыми поясами.
    start_date=pendulum.datetime(2025, 1, 1, tz="UTC"), 
    schedule=None,  # Устанавливаем `None` для ручного запуска DAG через UI
    tags=["branching", "example"], # Теги для фильтрации и поиска в интерфейсе Airflow
) as dag:
    
    # ------------------- ОПРЕДЕЛЕНИЕ ЗАДАЧ -------------------
    
    # Задача 1: `collect_order_data`
    # Используем PythonOperator для выполнения Python-функции
    def collect_order_data_func():
        """Имитирует сбор данных о новом заказе."""
        print("Собираю данные о новом заказе…")
        # Возвращаемое значение будет доступно в XComs для других задач
        return {"item_id": "SKU-001", "quantity": 10}

    collect_order_data = PythonOperator(
        task_id="collect_order_data", # Уникальный идентификатор задачи
        python_callable=collect_order_data_func # Функция, которая будет выполнена
    )

    # Задача 2: `check_stock`
    # Используем BranchPythonOperator для создания ветвления
    def check_stock_func(**context):
        """Проверяет наличие товара и выбирает следующую ветку DAG."""
        # Получаем объект Task Instance (ti) из контекста
        ti = context['ti']
        # Используем `xcom_pull` для получения данных, возвращённых предыдущей задачей (`collect_order_data`)
        order_info = ti.xcom_pull(task_ids='collect_order_data')
        quantity = order_info['quantity']
        
        # Логика ветвления:
        if quantity > 15:
            # Если товара достаточно, возвращаем ID задачи, которую нужно выполнить
            print("Товар есть в наличии. Продолжаю обработку.")
            return 'process_order'
        else:
            # Если товара нет, возвращаем ID задачи для уведомления
            print("Товара нет в наличии. Отправляю уведомление.")
            return 'notify_out_of_stock'

    check_stock = BranchPythonOperator(
        task_id="check_stock",
        python_callable=check_stock_func
    )
    
    # Задача 3a: `process_order` (ветка «наличие»)
    def process_order_func():
        """Имитирует обработку и отправку заказа."""
        print("Заказ успешно обработан и готов к отправке.")
        
    process_order = PythonOperator(
        task_id="process_order",
        python_callable=process_order_func
    )

    # Задача 3b: `notify_out_of_stock` (ветка «нет наличия»)
    def notify_out_of_stock_func():
        """Имитирует отправку уведомления клиенту."""
        print("Отправлено уведомление клиенту: 'Извините, товар временно отсутствует'.")
        
    notify_out_of_stock = PythonOperator(
        task_id="notify_out_of_stock",
        python_callable=notify_out_of_stock_func
    )

    # ------------------- ОПРЕДЕЛЕНИЕ ЗАВИСИМОСТЕЙ -------------------

    # Задаём последовательность выполнения задач с помощью оператора `>>`
    # Сначала выполнится `collect_order_data`, затем `check_stock`
    # После `check_stock` выполнится либо `process_order`, либо `notify_out_of_stock`
    # в зависимости от результата ветвления.
    # Задачи, которые не будут выполнены, получат статус "skipped" (пропущено)
    collect_order_data >> check_stock >> [process_order, notify_out_of_stock]

Пайплайн запускается вручную и состоит из четырёх задач, одна из которых определяет дальнейшую ветку выполнения.

Сначала задача collect_order_data имитирует получение информации о новом заказе. Она возвращает идентификатор товара и его количество. Airflow сохраняет эти данные в XCom — встроенном инструменте, с помощью которого задачи могут передавать друг другу информацию.

Затем запускается check_stock. Эта задача получает данные предыдущего этапа через XCom и проверяет количество товара. Для неё используется BranchPythonOperator, который позволяет выбрать одну из нескольких веток DAG.

В примере действует простое условие: если значение quantity больше 15, Airflow запускает задачу process_order. Она имитирует обработку заказа и подготовку к отправке. Если значение меньше 15, запускается notify_out_of_stock, которая имитирует отправку клиенту сообщения об отсутствии товара. Задача из второй, невыбранной ветки получает статус skipped и не выполняется.

Зависимости между задачами задаются оператором >>:

collect_order_data >> check_stock >> process_order / notify_out_of_stock

Таким образом в этом коде видно сразу несколько возможностей Airflow: создание DAG, выполнение Python-функций в отдельных задачах, передачу данных между ними через XCom и ветвление пайплайна в зависимости от результата проверки.

Чтобы проверить код, перезапустите Airflow. Для этого нажмите Ctrl + C в терминале и снова запустите команду:

airflow standalone

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

Структура DAG-а с несколькими действиями
Скриншот: Apache Airflow / Google Chrome / Skillbox Media

Что дальше

На этом базовое знакомство с Apache Airflow можно закончить. Мы разобрали, зачем нужен этот инструмент, как устроены DAG и основные компоненты системы, а также создали и запустили несколько простых пайплайнов.

Следующий шаг — перейти от учебных примеров к задачам, которые встречаются в реальных проектах. Например, можно научиться подключать Airflow к базам данных и внешним сервисам, передавать данные между задачами, настраивать повторные попытки и обработку ошибок. Отдельно стоит разобраться с мониторингом: как отслеживать состояние DAG, находить причины сбоев и настраивать уведомления.

Если планируете использовать Airflow в рабочем проекте, пригодятся и темы, связанные с инфраструктурой. Стоит изучить разные способы запуска задач, работу с Docker и Kubernetes, настройку базы метаданных, управление доступом и развёртывание Airflow в production-среде.

Лучший источник актуальной информации — официальная документация Apache Airflow. В ней подробно разобраны основные концепции, создание DAG, настройка задач, архитектура и развёртывание системы. Есть и отдельные пошаговые руководства, например по созданию первого рабочего процесса и использованию TaskFlow API. Документация регулярно обновляется вместе с Airflow, поэтому при возникновении конкретного технического вопроса лучше начинать именно с неё.

При выборе любых образовательных материалов, обращайте внимание на то, что актуальная стабильная ветка Airflow — 3.x, а поддержка Airflow 2 завершилась в апреле 2026 года. Поэтому старые курсы хорошо подходят для изучения базовых принципов, но конкретные инструкции лучше сверять с современной документацией.

Курс с трудоустройством: «Профессия Data scientist + ИИ» Узнать о курсе
Понравилась статья?
Да

Пользуясь нашим сайтом, вы соглашаетесь с тем, что мы используем cookies 🍪

Ссылка скопирована