Apache Airflow стал де-факто стандартом для оркестрации сложных рабочих процессов обработки данных. С ростом числа DAGs и их критической важности, ручной мониторинг через веб-интерфейс становится неэффективным и трудоемким. Возникает острая потребность в программном доступе к информации о состоянии выполнения DAGs и их отдельных запусков (DAG Runs).
Именно здесь на помощь приходит Airflow REST API. Он предоставляет мощный и гибкий интерфейс для взаимодействия с Airflow, позволяя автоматизировать мониторинг, интегрировать Airflow с внешними системами оповещения или BI-инструментами, а также создавать собственные панели управления. Понимание того, как эффективно использовать API для получения статуса выполнения, является ключевым навыком для любого инженера данных или DevOps-специалиста, работающего с Airflow.
В этом руководстве мы подробно рассмотрим, как использовать Airflow API для получения статуса DAG и DAG Run, предоставим практические примеры кода и объясним, как интерпретировать полученные результаты.
Введение в Airflow API и концепции статусов DAG
После того как мы осознали критическую важность программного мониторинга рабочих процессов в Apache Airflow, пришло время углубиться в инструментарий, который делает это возможным. Airflow предоставляет мощный REST API, позволяющий взаимодействовать с планировщиком и получать актуальную информацию о состоянии ваших DAG.
В этом разделе мы рассмотрим фундаментальные аспекты Airflow API, его назначение и ключевые концепции, такие как DAG и DAG Run, а также различные статусы, которые они могут принимать. Понимание этих основ является краеугольным камнем для эффективного использования API в целях мониторинга и автоматизации.
Что такое Apache Airflow API и зачем его использовать для мониторинга?
Apache Airflow API представляет собой программный интерфейс (Application Programming Interface), который позволяет взаимодействовать с Airflow на уровне кода, минуя веб-интерфейс или командную строку. Это RESTful API, предоставляющий стандартизированный набор эндпоинтов для выполнения различных операций, таких как управление DAGs, DAG Runs, задачами, переменными и подключениями.
Использование Airflow API для мониторинга критически важно по нескольким причинам:
-
Автоматизация: Позволяет автоматизировать проверку статусов выполнения DAGs и задач, что невозможно эффективно сделать вручную при большом количестве рабочих процессов.
-
Интеграция: Обеспечивает бесшовную интеграцию Airflow с внешними системами мониторинга, BI-инструментами или кастомными дашбордами, предоставляя актуальные данные о состоянии оркестрации.
-
Гибкость: Дает возможность создавать собственные скрипты и приложения для реагирования на определенные события (например, отправка уведомлений при сбое DAG) или для динамического управления рабочими процессами.
-
Масштабируемость: При росте числа DAGs и их запусков, программный доступ к информации становится единственным эффективным способом поддержания контроля и оперативного реагирования.
Основные понятия: DAG, DAG Run и их статусы выполнения
Для эффективного использования Airflow API критически важно понимать его фундаментальные концепции: DAG и DAG Run, а также их статусы выполнения.
-
DAG (Directed Acyclic Graph): Это основной строительный блок в Airflow. DAG представляет собой набор задач (тасков), организованных в определенной последовательности с заданными зависимостями. Он определяет что должно быть выполнено и в каком порядке, но сам по себе не является исполняемой сущностью. DAG – это, по сути, шаблон рабочего процесса.
-
DAG Run: Когда Airflow запускает DAG (вручную или по расписанию), создается экземпляр этого DAG, который называется DAG Run. Каждый DAG Run – это конкретное выполнение всех задач, определенных в DAG, для определенного момента времени или триггера. Именно DAG Run имеет конкретный статус выполнения.
Статусы выполнения:
В то время как сам DAG как шаблон не имеет динамического статуса (он либо активен, либо неактивен), каждый DAG Run проходит через ряд состояний, отражающих его прогресс и исход. Основные статусы DAG Run включают:
-
queued: Запуск DAG ожидает ресурсов для начала выполнения. -
running: Запуск DAG активен, и его задачи выполняются. -
success: Все задачи в DAG Run успешно завершены. -
failed: Одна или несколько задач в DAG Run завершились с ошибкой, и весь запуск считается неудачным. -
scheduled: Запуск был запланирован, но еще не начался. -
up_for_retry: Некоторые задачи в DAG Run находятся в состоянии повторной попытки.
Понимание этих различий между статичным определением DAG и динамическим экземпляром DAG Run со своим статусом является ключом к правильной интерпретации данных, получаемых через Airflow API.
Аутентификация и ключевые эндпоинты Airflow API
После того как мы разобрались с фундаментальными понятиями DAG, DAG Run и их статусами, следующим логичным шагом является понимание того, как получить эту информацию программно. Для эффективного взаимодействия с Airflow API и извлечения данных о статусах выполнения, необходимо прежде всего решить вопросы аутентификации. Без надлежащей авторизации доступ к ценным метрикам и состояниям будет невозможен.
В этом разделе мы рассмотрим различные методы аутентификации, которые Airflow API предлагает для безопасного доступа. Кроме того, мы сделаем обзор ключевых эндпоинтов REST API, которые специально предназначены для запроса информации о DAGs и их запусках, подготавливая почву для практических примеров.
Методы аутентификации для доступа к Airflow REST API
Для безопасного взаимодействия с Airflow REST API и получения статусов DAGs, необходима аутентификация. Airflow поддерживает несколько методов, обеспечивающих контролируемый доступ к его ресурсам:
-
Базовая аутентификация (Basic Authentication): Это наиболее распространенный и простой метод. Вы передаете имя пользователя и пароль, закодированные в Base64, в заголовке
Authorizationкаждого запроса. Пользователи и их роли управляются через веб-интерфейс Airflow или конфигурационные файлы. Убедитесь, что используемый пользователь имеет достаточные права для чтения информации о DAGs и DAG Runs. -
Токены API (API Tokens): Начиная с Airflow 2.0, появилась возможность генерировать токены API для конкретных пользователей. Эти токены представляют собой долгоживущие учетные данные, которые можно использовать вместо имени пользователя и пароля. Токены API более безопасны для использования в скриптах и автоматизированных системах, так как их можно отозвать в любой момент без изменения пароля пользователя. Токен передается в заголовке
AuthorizationкакBearer <ваш_токен>.
Выбор метода зависит от требований безопасности вашей инфраструктуры и удобства интеграции. Для автоматизированных скриптов предпочтительнее использовать токены API.
Обзор эндпоинтов для получения информации о DAGs и DAG Runs
После успешной аутентификации, о которой мы говорили ранее, можно приступать к взаимодействию с Airflow REST API для получения необходимой информации. Airflow предоставляет ряд интуитивно понятных эндпоинтов для работы с DAGs и их запусками (DAG Runs), которые являются ключевыми для мониторинга статусов.
Основные эндпоинты для получения информации о DAGs и DAG Runs:
-
/dags: Этот эндпоинт позволяет получить список всех DAGs, доступных в вашей инсталляции Airflow. В ответе содержится базовая информация о каждом DAG, включая егоdag_idи общее состояние (например,is_paused). -
/dags/{dag_id}: Используется для получения детальной информации о конкретном DAG по его идентификатору (dag_id). Здесь можно найти более подробные метаданные, но для статуса выполнения DAG Run потребуется другой эндпоинт. -
/dags/{dag_id}/dagRuns: Ключевой эндпоинт для мониторинга. Он возвращает список всех запусков (DAG Runs) для указанногоdag_id. Каждый элемент списка содержит важные данные, такие какdag_run_id,state(статус выполнения:success,failed,runningи т.д.),start_dateиend_date. -
/dags/{dag_id}/dagRuns/{dag_run_id}: Если вам нужна подробная информация о конкретном запуске DAG, этот эндпоинт предоставит все детали, включая его текущийstateи другие атрибуты, специфичные для данного запуска.
Все эти эндпоинты преимущественно используют HTTP-метод GET для извлечения данных. Понимание их структуры и возвращаемых данных является основой для эффективного программного мониторинга Airflow.
Практическое получение статуса DAG и DAG Run (с примерами кода)
После того как мы ознакомились с методами аутентификации и ключевыми эндпоинтами Airflow REST API, пришло время перейти от теории к практике. В этом разделе мы продемонстрируем, как программно взаимодействовать с Airflow API для получения статуса выполнения DAG и его отдельных запусков (DAG Runs), используя реальные примеры кода.
Мы рассмотрим пошаговые инструкции и готовые фрагменты кода на Python с использованием библиотеки requests, а также команды cURL. Это позволит вам не только понять принципы работы, но и сразу применить полученные знания для мониторинга ваших рабочих процессов.
Получение статуса конкретного DAG через Python (requests) и cURL
Для получения общего статуса конкретного DAG, включая его состояние (например, поставлен ли он на паузу) и информацию о последнем запуске, можно использовать эндпоинт /api/v1/dags/{dag_id}. Этот эндпоинт предоставляет метаданные о DAG, а также краткую информацию о его последнем DAG Run.
Пример с cURL
Для запроса статуса DAG с dag_id = example_dag:
curl -X GET "http://localhost:8080/api/v1/dags/example_dag" \
-H "Authorization: Basic <ВАШИ_УЧЕТНЫЕ_ДАННЫЕ_BASE64>"
В ответе вы получите JSON-объект, содержащий поля dag_id, is_paused, last_dag_run (с полем state для последнего запуска) и другие метаданные.
Пример с Python (библиотека requests)
import requests
import base64
AIRFLOW_URL = "http://localhost:8080"
DAG_ID = "example_dag"
USERNAME = "airflow"
PASSWORD = "airflow"
credentials = f"{USERNAME}:{PASSWORD}"
encoded_credentials = base64.b64encode(credentials.encode()).decode()
headers = {
"Authorization": f"Basic {encoded_credentials}",
"Content-Type": "application/json"
}
response = requests.get(f"{AIRFLOW_URL}/api/v1/dags/{DAG_ID}", headers=headers)
if response.status_code == 200:
dag_info = response.json()
print(f"Статус DAG '{DAG_ID}':")
print(f" ID: {dag_info.get('dag_id')}")
print(f" Пауза: {dag_info.get('is_paused')}")
last_run_state = dag_info.get('last_dag_run', {}).get('state')
print(f" Состояние последнего запуска: {last_run_state if last_run_state else 'N/A'}")
else:
print(f"Ошибка при получении статуса DAG: {response.status_code} - {response.text}")
Этот код извлекает is_paused (показывает, активен ли DAG) и state последнего DAG Run, что является ключевыми индикаторами общего состояния DAG.
Запрос статуса конкретного запуска DAG (DAG Run) и его задач
После того как мы научились получать общий статус DAG, часто возникает необходимость углубиться в детали конкретного выполнения. Для этого Airflow API предоставляет эндпоинты для работы с DAG Run и Task Instance.
Получение статуса конкретного DAG Run
Эндпоинт /api/v1/dags/{dag_id}/dagRuns/{dag_run_id} позволяет получить подробную информацию о конкретном запуске DAG. dag_run_id обычно формируется автоматически (например, scheduled__2026-03-28T00:00:00+00:00) или задается при ручном запуске (например, manual__my_custom_run).
Пример с cURL:
curl -X GET "http://localhost:8080/api/v1/dags/my_example_dag/dagRuns/manual__2026-03-28T10:00:00+00:00" \
-H "Authorization: Basic <ВАШИ_УЧЕТНЫЕ_ДАННЫЕ_В_BASE64>"
Пример на Python:
import requests
import base64
AIRFLOW_URL = "http://localhost:8080"
DAG_ID = "my_example_dag"
DAG_RUN_ID = "manual__2026-03-28T10:00:00+00:00"
USERNAME = "airflow"
PASSWORD = "airflow"
auth_string = f"{USERNAME}:{PASSWORD}".encode("ascii")
base64_auth = base64.b64encode(auth_string).decode("ascii")
headers = {"Authorization": f"Basic {base64_auth}"}
response = requests.get(f"{AIRFLOW_URL}/api/v1/dags/{DAG_ID}/dagRuns/{DAG_RUN_ID}", headers=headers)
if response.status_code == 200:
dag_run_data = response.json()
# print(f"Статус DAG Run '{DAG_RUN_ID}': {dag_run_data.get('state')}")
Получение статуса задач (Task Instances) в DAG Run
Для детального мониторинга выполнения отдельных задач внутри конкретного DAG Run используется эндпоинт /api/v1/dags/{dag_id}/dagRuns/{dag_run_id}/taskInstances. Он возвращает список всех экземпляров задач для указанного запуска.
Пример с cURL:
curl -X GET "http://localhost:8080/api/v1/dags/my_example_dag/dagRuns/manual__2026-03-28T10:00:00+00:00/taskInstances" \
-H "Authorization: Basic <ВАШИ_УЧЕТНЫЕ_ДАННЫЕ_В_BASE64>"
В ответе вы получите массив объектов Task Instance, каждый из которых будет содержать поля task_id, state, start_date, end_date и другие. Это позволяет точно отслеживать прогресс и выявлять проблемные задачи.
Интерпретация результатов и продвинутый мониторинг
После того как мы успешно освоили методы получения статусов DAG и DAG Run с помощью Airflow API, следующим критически важным шагом становится правильная интерпретация полученных данных. Недостаточно просто получить JSON-ответ; необходимо понимать, что означают различные состояния, как они соотносятся друг с другом и какие выводы можно сделать для эффективного мониторинга и отладки.
В этом разделе мы углубимся в анализ результатов запросов API, рассмотрим ключевые различия между статусами самого DAG и его конкретных запусков, а также обсудим, как использовать эту информацию для построения продвинутых систем отслеживания и оперативного реагирования на любые отклонения в работе ваших конвейеров данных.
Различия между статусом DAG и статусом DAG Run: что искать в ответе API
При работе с Airflow API важно четко различать понятия «статус DAG» и «статус DAG Run», поскольку они описывают разные аспекты состояния вашего рабочего процесса.
Статус DAG (сущности DAG)
Статус самого DAG как сущности в Airflow в основном отражает его конфигурационное состояние, а не результат выполнения. Ключевым параметром здесь является is_paused.
-
is_paused: true: Означает, что DAG приостановлен. Планировщик Airflow не будет создавать новые запуски (DAG Runs) для этого DAG, даже если его расписание наступило. Это полезно для временного отключения DAG без его удаления. -
is_paused: false: Означает, что DAG активен. Планировщик будет создавать новые DAG Runs в соответствии с заданным расписанием.
Этот статус можно получить из ответа эндпоинта /dags/{dag_id}. Он дает представление о том, будет ли DAG вообще запускаться.
Статус DAG Run (запуска DAG)
Статус DAG Run, напротив, описывает состояние конкретного экземпляра выполнения DAG. Каждый раз, когда DAG запускается (вручную или по расписанию), создается новый DAG Run со своим уникальным идентификатором и статусом. Основной параметр здесь — state.
Возможные значения state для DAG Run включают:
-
success: Все задачи в DAG Run успешно завершены. -
failed: Одна или несколько задач в DAG Run завершились с ошибкой. -
running: DAG Run в данный момент выполняется. -
queued: DAG Run ожидает ресурсов для запуска. -
scheduled: DAG Run был запланирован, но еще не начал выполняться. -
upstream_failed: Запуск не состоялся из-за сбоя в вышестоящем DAG (в случае зависимостей между DAG). -
skipped: Запуск был пропущен (например, из-за внешних условий).
Эти статусы доступны через эндпоинты /dags/{dag_id}/dagRuns и /dags/{dag_id}/dagRuns/{dag_run_id}. Они показывают, как прошел или проходит конкретный запуск.
Что искать в ответе API
При мониторинге важно проверять оба параметра:
-
is_paused: Убедитесь, что DAG не приостановлен, если вы ожидаете его выполнения. -
stateпоследнего DAG Run: Проверьте статус последнего или интересующего вас запуска, чтобы понять, успешно ли он завершился, все еще выполняется или произошел сбой.
Например, DAG может быть active (is_paused: false), но его последний DAG Run может иметь статус failed. Это указывает на проблему в самом рабочем процессе, а не в его планировании. И наоборот, если DAG paused (is_paused: true), то отсутствие новых running или queued DAG Runs является ожидаемым поведением.
Отслеживание и интерпретация различных состояний выполнения, включая получение логов
После того как мы разобрались с различиями между статусом DAG и статусом DAG Run, перейдем к более глубокой интерпретации состояний выполнения и методам получения детальной информации, такой как логи.
Интерпретация состояний DAG Run
Статус state объекта DAG Run может принимать различные значения, каждое из которых указывает на определенный этап жизненного цикла выполнения:
-
scheduled: Запуск DAG запланирован, но еще не поставлен в очередь на выполнение. -
queued: Запуск DAG находится в очереди и ожидает доступности исполнителя. -
running: DAG Run активно выполняется. -
success: Все задачи в DAG Run успешно завершены. -
failed: Одна или несколько задач в DAG Run завершились с ошибкой, и DAG Run в целом считается неудачным. -
up_for_retry: Задача завершилась с ошибкой, но настроены повторные попытки, и она будет перезапущена. -
skipped: Задача была пропущена, например, из-за условий ветвления.
Получение логов для детального анализа
Для глубокого анализа причин сбоев или аномального поведения критически важно иметь доступ к логам выполнения задач. Airflow API предоставляет эндпоинты для получения логов конкретных экземпляров задач (Task Instances). Используя эндпоинт /dags/{dag_id}/dagRuns/{dag_run_id}/taskInstances/{task_id}/logs/{log_id}, вы можете получить логи для определенной задачи в конкретном запуске DAG. Это позволяет программно извлекать информацию, которая обычно доступна через веб-интерфейс, и интегрировать ее в собственные системы мониторинга или отчетности. Например, запрос к этому эндпоинту вернет текстовое содержимое лога, которое можно парсить для выявления ошибок или ключевых событий.
Заключение
Мы рассмотрели, как Airflow API предоставляет мощный инструментарий для программного мониторинга и управления статусами DAG и DAG Run. Использование REST API позволяет не только оперативно получать информацию о ходе выполнения ваших пайплайнов, но и интегрировать Airflow в существующие системы мониторинга и автоматизации. Это критически важно для поддержания стабильности и эффективности сложных рабочих процессов, обеспечивая прозрачность и контроль над оркестрацией данных.