Первый commit
@@ -0,0 +1,80 @@
|
||||
# Почему Apache Airflow стал незаменимым инструментом для работы с данными
|
||||
# Почему Apache Airflow стал незаменимым инструментом для работы с данными
|
||||
|
||||
Обычно работа с автоматизацией процессов обработки информации начинается с ручного управления задачами. Например, в машинном обучении это может включать подготовку наборов данных, обучение моделей, анализ результатов и развертывание решений в рабочей среде. По мере роста команды и развития продукта эти процессы усложняются: увеличивается количество повторяющихся операций, появляются зависимости между задачами, и каждая из них приобретает всё большее значение для бизнеса. В результате формируется полноценный конвейер задач, требующий регулярного запуска.
|
||||
|
||||

|
||||
|
||||
Аналогичная ситуация возникает при обработке данных: в определенный момент необходимо собрать актуальную информацию, преобразовать её и выполнить различные операции — создать витрину данных и сохранить в базу, обучить модель машинного обучения, подготовить отчет в Excel и разослать его по электронной почте. Вариантов множество, и для решения таких задач требуются специализированные инструменты.
|
||||
|
||||
Одним из таких решений является Apache Airflow. Благодаря поддержке сообщества разработчиков он стал стандартом де-факто в своей области.
|
||||
|
||||
В сфере работы с данными Airflow часто называют инструментом для пакетной обработки данных в рамках подхода ETL (Extract, Transform, Load). Однако важно понимать, что Airflow не является классической ETL-системой, а скорее помогает организовывать весь процесс извлечения, преобразования и загрузки данных через Python-скрипты.
|
||||
|
||||
Наиболее точное определение Airflow — это оркестратор рабочих процессов. Задачи и их зависимости описываются в скриптах на Python, а Airflow отвечает за планирование и выполнение кода на основе заданных конфигураций.
|
||||
|
||||
# Преимущества использования Apache Airflow
|
||||
|
||||
Вот ключевые причины, по которым стоит выбрать Airflow:
|
||||
|
||||
- **Открытый исходный код**. Изначально разработанный как внутренний проект Airbnb, Airflow нуждался в активном сообществе для своего развития. Именно поэтому он был выпущен как open-source решение. Сегодня Airflow поддерживается и управляется как один из флагманских проектов Apache Software Foundation.
|
||||
|
||||
💡 В дальнейших материалах вы можете встретить как "Apache Airflow", так и просто "Airflow" — речь идет об одном и том же инструменте с открытым исходным кодом, разрабатываемом под эгидой Apache.
|
||||
|
||||
- **Интуитивный веб-интерфейс**. Пользовательский интерфейс позволяет легко отслеживать все рабочие процессы, изменять их параметры, запускать или останавливать выполнение. Хотя работа возможна и через командную строку, веб-интерфейс значительно понижает порог входа в технологию, делая Airflow доступным не только для инженеров данных, но и для аналитиков, разработчиков, системных администраторов и DevOps-специалистов.
|
||||
|
||||
- **Основан на Python**. Вся конфигурация создается с использованием языка Python, включая настройку расписаний и скриптов для выполнения задач. Это позволяет использовать привычные Python-библиотеки и классы для создания рабочих процессов, избавляя от необходимости работать с JSON или XML конфигурационными файлами. Кроме того, Python является стандартом де-факто для специалистов в области Big Data и Data Science.
|
||||
|
||||
- **Широкое распространение**. Airflow активно используется такими компаниями, как Airbnb, Intel, PayPal, WePay и Yahoo!, став незаменимым инструментом в арсенале инженера данных. Опыт работы с Airflow часто указывается как требование в вакансиях на позицию Data Engineer.
|
||||
|
||||
- **Простота внедрения**. Легкая установка, быстрый старт, удобная визуализация, возможность автоматического создания большого количества задач и широкие возможности кастомизации.
|
||||
|
||||
- **Масштабируемость**. Модульная архитектура и поддержка очередей сообщений (Celery/Dask) позволяют работать с неограниченным количеством DAG.
|
||||
|
||||
- **Встроенный репозиторий метаданных**. На базе библиотеки SQLAlchemy хранятся состояния задач, DAG, глобальные переменные и другая служебная информация.
|
||||
|
||||
- **Богатая экосистема интеграций**. Поддержка различных баз данных (MySQL, PostgreSQL, DynamoDB, Hive), хранилищ Big Data (HDFS, Amazon S3) и облачных платформ (Google Cloud Platform, Amazon Web Services, Microsoft Azure).
|
||||
|
||||
- **Расширяемый REST API**. Возможность легко интегрировать Airflow в существующую IT-инфраструктуру и гибко настраивать конвейеры данных, например, передавая параметры в DAG через HTTP-запросы.
|
||||
|
||||
# Когда Airflow может не подойти
|
||||
|
||||
Существуют сценарии, для которых этот инструмент не является оптимальным выбором:
|
||||
|
||||
- **Потоковая обработка данных**. Airflow изначально разработан для выполнения повторяющихся пакетных задач, а не для обработки потоковых данных, где события могут быть распределены во времени. В таких случаях лучше рассмотреть специализированные решения.
|
||||
|
||||
- **Требования к знанию Python**. Поскольку Airflow полностью основан на Python, работа с ним предполагает хорошее владение этим языком программирования. Специалисты, знакомые только с SQL или другими языками, могут испытывать трудности. В таких ситуациях можно рассмотреть решения с минимальным количеством кода и упором на графический интерфейс (например, Azure Data Factory) или упростить задачу до статического процесса с запуском по расписанию.
|
||||
|
||||
- **Сложность поддержки чистоты кода**. При работе с конвейерами обработки данных (представленными в Airflow как DAG), код может быстро усложниться и стать непонятным для других разработчиков. Поэтому успешная работа с Airflow требует от инженеров данных хорошего владения Python и соблюдения принципов проектирования, таких как KISS, правильного именования переменных и написания документации.
|
||||
|
||||
Если ваш сценарий не включает потоковую обработку данных (например, анализ событий в мобильном приложении) и вы не испытываете сложностей с Python, то Airflow, скорее всего, станет отличным выбором, а его базовый функционал покроет большинство ваших потребностей.
|
||||
|
||||
# Практическое применение Airflow
|
||||
|
||||
На практике Apache Airflow используется в следующих сценариях:
|
||||
|
||||
- Интеграция данных из множества информационных систем (внутренние и внешние базы данных, файловые хранилища, облачные приложения и т.д.)
|
||||
- Загрузка информации в корпоративное озеро данных (Data Lake)
|
||||
- Создание уникальных конвейеров доставки и обработки больших объемов данных (data pipeline)
|
||||
- Управление конфигурацией конвейеров данных как кодом в соответствии с DevOps-подходом
|
||||
- Автоматизация разработки, планирования и мониторинга пакетных процессов обработки данных
|
||||
|
||||
💡 Пакетная обработка данных (batch processing) подразумевает обработку информации крупными порциями, где у каждой задачи есть четко определенное начало и конец.
|
||||
|
||||
# Как компании развертывают Airflow
|
||||
|
||||
В российских компаниях Airflow часто развертывается на собственных серверах силами DevOps-инженеров или системных администраторов. Это связано с требованиями законодательства о хранении и обработке пользовательских данных на территории РФ. Такой подход называют on-premise (или on-prem), что означает развертывание и администрирование инструмента собственными силами на арендованных или приобретенных мощностях в дата-центрах.
|
||||
|
||||
Одновременно с этим крупные облачные провайдеры предлагают Airflow в виде управляемого сервиса:
|
||||
|
||||
- **Astronomer** создал SaaS-платформу вокруг Airflow с расширенными возможностями мониторинга, оповещения об инцидентах и DevOps-инструментами, развертываемыми на кластерах Kubernetes.
|
||||
- **Cloud Composer** — управляемый облачный сервис Airflow от Google Cloud Platform (GCP), тесно интегрированный с другими сервисами GCP.
|
||||
- **Amazon Managed Workflows for Apache Airflow (MWAA)** — аналогичное решение от Amazon Web Services (AWS).
|
||||
|
||||
Такой "облачный" подход часто предоставляется как managed-сервис, что означает отсутствие необходимости заботиться о развертывании и администрировании — этим занимается провайдер. Однако у такого подхода есть и недостатки, в первую очередь высокая стоимость при больших объемах данных или интенсивном использовании.
|
||||
|
||||
Облачное развертывание Airflow преимущественно распространено в зарубежных компаниях и стартапах.
|
||||
|
||||
В этом материале вы познакомились с предпосылками появления инструментов управления процессами, подобных Airflow. Вы узнали, в каких сценариях Airflow наиболее эффективен, а в каких случаях стоит рассмотреть альтернативные решения, а также как этот инструмент применяется на практике в крупных российских и международных компаниях.
|
||||
|
||||
💡 Обратите внимание, что в данном модуле рассматривается Airflow версии 2.5.
|
||||
@@ -0,0 +1,101 @@
|
||||
# Понятное введение в ключевые понятия Airflow
|
||||
|
||||
В этом материале мы разберем фундаментальные концепции Apache Airflow, которые вам обязательно понадобятся при создании пайплайнов обработки данных. Знание этих базовых терминов и компонентов поможет вам эффективно работать с системой.
|
||||
|
||||
# **Что такое DAG и почему он важен**
|
||||
|
||||
Сердце Apache Airflow — это DAG (Directed Acyclic Graph), что переводится как "направленный ациклический граф". Давайте разберем это сложное название по частям:
|
||||
|
||||
- **Граф** — это структура, где элементы (узлы) связаны между собой стрелками (ребрами)
|
||||
- **Направленный** — означает, что связи имеют четкое направление: от одного элемента к другому, как последовательные этапы в обработке данных
|
||||
- **Ациклический** — гарантирует, что вы не можете вернуться к уже пройденному узлу, избегая бесконечных циклов
|
||||
|
||||
Проще говоря, DAG — это упорядоченный набор задач, которые выполняются строго по расписанию и никогда не повторяются в рамках одного запуска. Создавая пайплайны для обработки данных, вы фактически создаете DAG.
|
||||
|
||||

|
||||
|
||||
# Шаги вашего пайплайна: Задачи
|
||||
|
||||
Задачи — это отдельные шаги в вашем процессе обработки данных. Каждая задача представляет собой конкретную бизнес-операцию:
|
||||
- Проверка наличия файла в хранилище
|
||||
- Выполнение SQL-запроса
|
||||
- Запуск Python-функции с бизнес-логикой
|
||||
- Отправка уведомления
|
||||
|
||||
💡 **Совет по проектированию**: Разбивайте сложные процессы на мелкие, атомарные задачи. Такой подход делает пайплайн более гибким, упрощает тестирование и помогает быстрее находить и исправлять ошибки.
|
||||
|
||||
На примере ниже видно, как задача E зависит от успешного завершения всех предыдущих задач:
|
||||
|
||||

|
||||
|
||||
Airflow позволяет создавать сложные сценарии:
|
||||
- Зависимости между разными DAG-ами с помощью TriggerDagRunOperator и ExternalTaskSensor
|
||||
- Условное выполнение задач в зависимости от результатов предыдущих шагов
|
||||
- Сложные ветвления и параллельные ветки выполнения
|
||||
|
||||

|
||||
|
||||
# Инструменты для выполнения: Операторы
|
||||
|
||||
Операторы — это готовые шаблоны, которые определяют, КАК выполнять ваши задачи. Думайте о них как о строительных блоках для вашего пайплайна.
|
||||
|
||||
Пример простого DAG с двумя задачами:
|
||||
|
||||
В этом примере создается DAG с идентификатором 'greeting_dag', который запускается каждые 10 минут. DAG включает в себя две задачи: 'start_process' (использует DummyOperator как точку старта) и 'greeting_task' (использует PythonOperator для выполнения функции приветствия).
|
||||
|
||||
```python
|
||||
from airflow import DAG
|
||||
from airflow.operators.dummy import DummyOperator
|
||||
from airflow.operators.python import PythonOperator
|
||||
from datetime import datetime
|
||||
|
||||
def greet_user():
|
||||
print('Welcome to Airflow!')
|
||||
|
||||
dag = DAG('greeting_dag',
|
||||
description='Simple Greeting DAG',
|
||||
schedule_interval='*/10 * * * *',
|
||||
start_date=datetime(2023, 1, 1),
|
||||
catchup=False)
|
||||
|
||||
start_task = DummyOperator(task_id='start_process', retries=2)
|
||||
greeting_task = PythonOperator(task_id='greeting_task', python_callable=greet_user)
|
||||
|
||||
start_task >> greeting_task
|
||||
```
|
||||
|
||||
В этом примере задача `greeting_task` использует `PythonOperator` для запуска функции `greet_user`, которая выводит приветственное сообщение в логи.
|
||||
|
||||
Airflow предоставляет множество готовых операторов:
|
||||
- **PythonOperator** — выполнение Python-кода
|
||||
- **BashOperator** — запуск команд и скриптов в терминале
|
||||
- **PostgresOperator** — выполнение SQL-запросов (аналоги есть для MySQL, Oracle, Hive)
|
||||
- **EmailOperator** — отправка электронных писем
|
||||
- **DummyOperator** — "заглушка" для организации структуры пайплайна
|
||||
|
||||
💡 **Важное различие**: Задача определяет ЧТО нужно сделать (бизнес-логика), а оператор — КАК это сделать (техническая реализация).
|
||||
|
||||
# Ожидание событий: Сенсоры
|
||||
|
||||
Сенсоры — это специальный тип операторов, предназначенный для ожидания определенных условий или событий. Они идеально подходят для событийно-ориентированных пайплайнов.
|
||||
|
||||
Популярные сенсоры включают:
|
||||
- **PythonSensor** — ожидает, пока функция не вернет `True`
|
||||
- **S3Sensor** — проверяет наличие файла в S3-бакете
|
||||
- **RedisPubSubSensor** — ждет поступления сообщения в очередь
|
||||
- **RedisKeySensor** — проверяет существование ключа в Redis
|
||||
|
||||
Airflow поддерживает расширение функциональности через providers и позволяет создавать собственные операторы и сенсоры под специфические нужды.
|
||||
|
||||
# Краткое резюме
|
||||
|
||||
Давайте закрепим ключевые понятия:
|
||||
|
||||
- **DAG** — ваш пайплайн обработки данных, объединяющий задачи в логическую последовательность
|
||||
- **Задача** — отдельный шаг в пайплайне, представляющий конкретную бизнес-операцию
|
||||
- **Оператор** — технический инструмент для выполнения задачи (Python, SQL, Bash и т.д.)
|
||||
- **Сенсор** — специальный оператор для ожидания внешних событий или условий
|
||||
|
||||
Взаимосвязь этих компонентов можно представить так:
|
||||
|
||||

|
||||
@@ -0,0 +1,59 @@
|
||||
# Как устроен Airflow внутри
|
||||
|
||||
В этом уроке мы разберём, как устроен Apache Airflow изнутри. Вы узнаете, из каких компонентов состоит эта система и как они взаимодействуют между собой для выполнения задач по обработке данных. В конце мы проследим полный путь выполнения задачи от начала до конца.
|
||||
|
||||
## Основные компоненты Airflow
|
||||
|
||||
Apache Airflow состоит из нескольких взаимосвязанных компонентов, которые совместно обеспечивают выполнение задач, отслеживание их статусов, перезапуск при ошибках и логирование операций.
|
||||
|
||||
На схеме ниже показано, как эти компоненты взаимодействуют между собой:
|
||||
|
||||

|
||||
*Как устроена система Airflow*
|
||||
|
||||
Давайте подробно рассмотрим каждый компонент системы.
|
||||
|
||||
### Файлы DAG
|
||||
Это Python-файлы, в которых инженеры данных описывают последовательности операций по обработке данных. Все эти файлы хранятся в специальной директории, которую Airflow постоянно отслеживает на наличие изменений.
|
||||
|
||||
### Конфигурационный файл (airflow.cfg)
|
||||
Главный файл настроек Airflow, где задаются все параметры работы системы: частота проверки расписаний, пути к файлам, настройки подключения к базе данных и многое другое. Практически все компоненты Airflow обращаются к этому файлу для получения необходимых параметров.
|
||||
|
||||
### База данных метаданных
|
||||
Это хранилище служебной информации о работе системы: статусы задач, время их запуска и завершения, зависимости между задачами и другая техническая информация. В качестве базы данных обычно используют PostgreSQL или MySQL.
|
||||
|
||||
### Веб-интерфейс (Web Server)
|
||||
Графический интерфейс, с которым взаимодействуют пользователи Airflow. Он позволяет визуализировать DAG-и, отслеживать статусы задач, просматривать логи и управлять выполнением пайплайнов. Веб-интерфейс построен на Python с использованием фреймворка Flask.
|
||||
|
||||
### Планировщик (Scheduler)
|
||||
Центральный компонент Airflow, который отвечает за координацию всех процессов. Планировщик отслеживает расписания запуска DAG-ов, определяет, какие задачи нужно выполнять, и управляет их статусами.
|
||||
|
||||
### Исполнитель (Executor)
|
||||
Компонент, который работает в тесной связке с планировщиком и определяет, как и где будут выполняться задачи. В зависимости от инфраструктуры можно использовать разные типы исполнителей: LocalExecutor (для локального выполнения), CeleryExecutor, KubernetesExecutor и другие.
|
||||
|
||||
### Рабочие процессы (Workers)
|
||||
Это "исполнители" задач, которые фактически выполняют код задач. Рабочие процессы сообщают о своём статусе планировщику и исполнителю. В распределённых системах рабочие процессы могут работать на отдельных серверах для повышения надёжности и производительности.
|
||||
|
||||
## Как всё работает вместе
|
||||
|
||||
Теперь давайте проследим, как происходит выполнение задачи в Airflow от начала до конца.
|
||||
|
||||
Представим, что инженер данных создал DAG с ежедневным расписанием на 12 часов дня и поместил файл в директорию DAG-ов.
|
||||
|
||||
1. Планировщик постоянно проверяет текущее время и сравнивает его с расписаниями всех DAG-ов.
|
||||
|
||||
2. В 12 часов планировщик обнаруживает, что настало время запустить наш DAG. Он считывает файл DAG, анализирует зависимости между задачами и определяет порядок их выполнения.
|
||||
|
||||
3. Информация о задачах и их начальных статусах записывается в базу данных метаданных.
|
||||
|
||||
4. Планировщик обращается к исполнителю, указанному в конфигурационном файле, чтобы тот выделил ресурсы для выполнения задач.
|
||||
|
||||
5. Исполнитель запускает рабочие процессы для выполнения задач. В случае LocalExecutor задачи выполняются как локальные подпроцессы, а при использовании других исполнителей — на удалённых рабочих узлах.
|
||||
|
||||
6. По мере выполнения задач рабочие процессы обновляют их статусы в базе данных метаданных, что позволяет планировщику отслеживать прогресс и запускать следующие задачи в зависимости от успешности предыдущих.
|
||||
|
||||
Таким образом, Airflow обеспечивает надёжное и контролируемое выполнение сложных пайплайнов обработки данных благодаря чёткому разделению ответственности между компонентами системы.
|
||||
|
||||
Для более глубокого изучения основных концепций Airflow рекомендуем обратиться к [официальной документации](https://airflow.apache.org/docs/apache-airflow/stable/concepts/overview.html).
|
||||
|
||||
В этом уроке вы узнали, как устроена система Airflow изнутри, познакомились с её ключевыми компонентами и поняли, как они взаимодействуют для выполнения задач по обработке данных.
|
||||
@@ -0,0 +1,114 @@
|
||||
# Знакомство с интерфейсом Airflow для начинающих
|
||||
|
||||
Когда вы создаете ETL-процессы для автоматической обработки данных, важно иметь удобный способ отслеживать их работу, находить ошибки и управлять выполнением. Именно для этого в Apache Airflow предусмотрен веб-интерфейс — ваш главный помощник в повседневной работе с пайплайнами.
|
||||
|
||||
В этом материале мы подробно разберем, как устроен интерфейс Airflow версии 2.5 и какие возможности он предоставляет для мониторинга и управления вашими процессами обработки данных.
|
||||
|
||||
# Домашняя страница Airflow
|
||||
|
||||
После входа в систему вы попадете на главную страницу — центральную панель управления, где собрана вся ключевая информация о ваших пайплайнах. Не пугайтесь обилия элементов и цветов — все устроено логично и интуитивно понятно.
|
||||
|
||||
По сути, это обычная таблица, где каждая строка представляет собой один DAG (Directed Acyclic Graph) — ваш пайплайн обработки данных, а столбцы содержат различную информацию о нем.
|
||||
|
||||
## Основные элементы главной страницы
|
||||
|
||||
**Название DAG** — в первом столбце отображается список всех зарегистрированных в системе пайплайнов. По умолчанию Airflow включает демонстрационные примеры различных операторов. Список отсортирован по алфавиту для удобства поиска.
|
||||
|
||||

|
||||
|
||||
**Переключатели активности** — напротив каждого DAG находится кнопка-выключатель, позволяющая мгновенно активировать или деактивировать пайплайн прямо из веб-интерфейса без изменения кода.
|
||||
|
||||

|
||||
|
||||
**Владелец процесса** — каждый пайплайн имеет ответственного владельца. Это особенно полезно в командной работе, когда несколько инженеров создают и поддерживают различные ETL-процессы. Владелец отвечает за мониторинг и корректную работу своего DAG.
|
||||
|
||||

|
||||
|
||||
**Статус выполнения** — цветные индикаторы с цифрами показывают количество и состояние последних запусков DAG:
|
||||
- 🔴 Красный — завершено с ошибкой (failed)
|
||||
- 🟡 Желтый — ожидает повторного запуска (retry)
|
||||
- 🟢 Зеленый — выполняется в данный момент (running)
|
||||
- 🟢 Тёмно-зеленый — успешно завершено (success)
|
||||
|
||||

|
||||
|
||||
**Расписание** — указывает, когда и с какой периодичностью запускается пайплайн. Используется формат cron, который может показаться сложным на первый взгляд. Для перевода cron-выражений в понятный формат рекомендуем использовать сервис [Crontab.guru](https://crontab.guru/).
|
||||
|
||||

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

|
||||
|
||||
**Статус задач** — детальная информация о последнем запуске: сколько задач находится в каждом статусе. Это помогает быстро оценить общее состояние пайплайна без необходимости погружаться в детали.
|
||||
|
||||

|
||||
|
||||
**Быстрые действия** — в последнем столбце расположены кнопки для немедленного выполнения операций: запуск, обновление и удаление DAG. На практике этими кнопками пользуются редко.
|
||||
|
||||
Главная страница дает вам общее представление о состоянии всех ваших процессов. Но для детальной работы с конкретным пайплайном нужно перейти внутрь — просто кликните по названию интересующего DAG.
|
||||
|
||||

|
||||
|
||||
# Страница конкретного DAG
|
||||
|
||||
После перехода внутрь DAG вы увидите набор вкладок с различной информацией: от визуального представления структуры пайплайна до детальных логов выполнения и исходного кода.
|
||||
|
||||
## Древовидное представление (Tree View)
|
||||
|
||||
По умолчанию открывается вкладка с древовидной структурой задач. Здесь отображаются все запуски DAG с указанием статуса каждой задачи, времени выполнения и других метрик мониторинга.
|
||||
|
||||
Вы можете увидеть:
|
||||
- Состав DAG и последовательность выполнения задач
|
||||
- Тип оператора для каждой задачи
|
||||
- Историю запусков в виде цветных квадратов напротив каждой задачи
|
||||
|
||||

|
||||

|
||||
|
||||
## Графическое представление (Graph View)
|
||||
|
||||
Когда DAG содержит много задач, древовидное представление может быть неудобным. В таких случаях используйте вкладку Graph View — она показывает пайплайн в виде наглядного графа с четкими связями между задачами.
|
||||
|
||||

|
||||

|
||||
|
||||
При клике на любую задачу открывается подробное окно с двумя основными разделами:
|
||||
|
||||
### Информация о задаче
|
||||
- **Просмотр логов** — переход к странице с полным выводом выполнения задачи (одна из самых часто используемых функций)
|
||||
- **Детали выполнения** — подробная информация о конкретном запуске задачи
|
||||
|
||||
### Управление задачей
|
||||
Доступны четыре основных действия:
|
||||
- **Запустить** — выполнить задачу немедленно
|
||||
- **Очистить состояние** — сбросить статус задачи для повторного выполнения
|
||||
- **Отметить как неудачную** — вручную установить статус ошибки
|
||||
- **Отметить как успешную** — вручную установить статус успеха
|
||||
|
||||
Каждое действие можно комбинировать с дополнительными опциями:
|
||||
- **Игнорировать зависимости** — запуск без проверки зависимостей от других задач
|
||||
- **Работать с прошлыми/будущими запусками** — применить действие ко всем запускам в определенном временном диапазоне
|
||||
- **Влиять на связанные задачи** — применить действие к предыдущим (upstream) или последующим (downstream) задачам
|
||||
|
||||
Наиболее популярная комбинация — **Downstream + Recursive + Clear**, которая сбрасывает текущую задачу и все зависящие от нее задачи в рамках одного запуска.
|
||||
|
||||
## Анализ времени выполнения (Task Duration)
|
||||
|
||||
Вкладка Task Duration автоматически строит графики на основе истории выполнения, показывая, сколько времени занимает каждая задача при каждом запуске. Это помогает выявлять узкие места и отслеживать изменения производительности.
|
||||
|
||||

|
||||
|
||||
## Диаграмма Ганта (Gantt)
|
||||
|
||||
Диаграмма Ганта визуализирует распределение времени выполнения задач в рамках одного запуска DAG. Это отличный инструмент для определения самых ресурсоемких операций и планирования оптимизации.
|
||||
|
||||

|
||||
|
||||
## Исходный код (Code)
|
||||
|
||||
Вкладка Code отображает актуальный код DAG, который Airflow использует для выполнения. Это особенно полезно для проверки, что изменения из вашего Git-репозитория успешно загружены в систему и готовы к выполнению.
|
||||
|
||||

|
||||
|
||||
Теперь вы знакомы с основными возможностями веб-интерфейса Airflow для мониторинга и управления вашими процессами обработки данных. Эти знания помогут вам эффективно работать с пайплайнами и быстро решать возникающие проблемы.
|
||||
@@ -0,0 +1,286 @@
|
||||
# Основы построения DAG-файлов в Airflow
|
||||
|
||||
В предыдущем занятии вы познакомились с пользовательским интерфейсом Airflow и, вероятно, задались вопросами: как создать свой собственный DAG? Где писать код? Какие именно инструкции использовать? В этом уроке вы получите ответы на все эти вопросы, создадите свой первый DAG, изучите базовую структуру кода и сможете загрузить его в систему для запуска и наблюдения за результатами.
|
||||
|
||||
# Как устроен код DAG-файла
|
||||
|
||||
Как вы уже знаете из предыдущих уроков, DAG представляет собой граф вычислений, состоящий из отдельных задач (tasks). На уровне кода DAG — это обычный Python-файл, который описывает все задачи в рамках пайплайна и определяет последовательность их выполнения.
|
||||
|
||||
Любой DAG-файл состоит из нескольких ключевых компонентов:
|
||||
- Импорт необходимых модулей и библиотек
|
||||
- Настройка параметров и инициализация объекта DAG
|
||||
- Создание отдельных задач с помощью операторов
|
||||
- Определение порядка выполнения задач
|
||||
|
||||
Порядок этих компонентов имеет значение, поскольку Airflow — это Python-библиотека, и обращение к еще не инициализированным объектам приведет к ошибкам выполнения.
|
||||
|
||||
Давайте рассмотрим простой пример DAG-файла и разберем его по частям.
|
||||
|
||||

|
||||
*Пример структуры DAG-файла*
|
||||
|
||||
## Импорт необходимых модулей
|
||||
|
||||
Все начинается с подключения нужных библиотек и классов. Обязательно импортируйте класс DAG для создания графа вычислений. Затем подключайте операторы, которые понадобятся для описания ваших задач. Также могут потребоваться дополнительные функции или объекты для работы с данными.
|
||||
|
||||
```python
|
||||
from datetime import datetime, timedelta
|
||||
|
||||
from airflow import DAG
|
||||
from airflow.operators.bash import BashOperator
|
||||
from airflow.utils.dates import days_ago
|
||||
```
|
||||
|
||||
## Настройка параметров DAG
|
||||
|
||||
Следующим шагом идет объявление DAG с необходимыми параметрами:
|
||||
|
||||
В примере ниже создается DAG с идентификатором "sample_dag", который будет использовать общие параметры, определенные в словаре `default_args`. В данном случае, владелец DAG - это 'data_team', а дата начала выполнения - 1 января 2021 года.
|
||||
|
||||
```python
|
||||
default_args = {
|
||||
'start_date': datetime(2021, 1, 1),
|
||||
'owner': 'data_team'
|
||||
}
|
||||
|
||||
dag = DAG(
|
||||
"sample_dag",
|
||||
default_args=default_args,
|
||||
schedule_interval=None,
|
||||
)
|
||||
```
|
||||
|
||||
Обратите внимание на словарь `default_args`. Это стандартная практика для хранения общих параметров, которые будут применяться ко всем задачам в DAG. Такой подход помогает избежать дублирования кода и делает его более читаемым. В примере мы указываем владельца пайплайна и дату начала его работы.
|
||||
|
||||
Ключевыми обязательными параметрами являются `owner` и `start_date`. Без них DAG не сможет быть корректно инициализирован и не появится в пользовательском интерфейсе Airflow. Остальные параметры являются опциональными и добавляются по мере необходимости.
|
||||
|
||||
Для изучения всех доступных параметров рекомендуем обращаться к официальной документации Airflow.
|
||||
|
||||
## Создание задач с помощью операторов
|
||||
|
||||
После настройки DAG следует этап создания задач. Для каждого шага обработки данных создается отдельная переменная с соответствующим оператором.
|
||||
|
||||
В Airflow существует множество готовых операторов для различных задач:
|
||||
- `BashOperator` — для выполнения bash-команд
|
||||
- `PythonOperator` — для запуска Python-функций
|
||||
- Специализированные операторы для работы с популярными системами обработки данных (например, Apache Spark)
|
||||
|
||||
Каждая задача должна иметь уникальный идентификатор (`task_id`), который используется Airflow для отображения задачи в интерфейсе.
|
||||
|
||||
Пример создания простых задач:
|
||||
|
||||
В этом примере мы создаем две задачи с использованием BashOperator. Первая задача с идентификатором 'show_time' выводит текущее время, а вторая задача с идентификатором 'process_data' имитирует обработку данных с задержкой 3 секунды и возможностью повторного запуска при ошибках (до 2 раз).
|
||||
|
||||
```python
|
||||
# Задача для вывода текущего времени
|
||||
t1 = BashOperator(
|
||||
task_id='show_time',
|
||||
bash_command='date',
|
||||
dag=dag
|
||||
)
|
||||
|
||||
# Задача для имитации обработки с возможностью повторных попыток
|
||||
t2 = BashOperator(
|
||||
task_id='process_data',
|
||||
bash_command='sleep 3',
|
||||
retries=2,
|
||||
dag=dag
|
||||
)
|
||||
```
|
||||
|
||||
## Определение последовательности выполнения
|
||||
|
||||
Завершающий этап — указание порядка выполнения задач. В Airflow для этого используются стрелочные операторы:
|
||||
|
||||
В приведенном примере задача с идентификатором 'process_data' будет запускаться только после успешного завершения задачи 'show_time'.
|
||||
|
||||
```python
|
||||
t1 >> t2 # t2 запускается после завершения t1
|
||||
```
|
||||
|
||||
Альтернативный синтаксис:
|
||||
|
||||
Этот синтаксис эквивалентен предыдущему примеру, просто записан в обратном порядке. Задача 'show_time' должна завершиться перед запуском задачи 'process_data'.
|
||||
```python
|
||||
t2 << t1 # t1 должна завершиться перед запуском t2
|
||||
```
|
||||
|
||||
Можно также группировать задачи в списки для создания более сложных зависимостей:
|
||||
|
||||
В этом примере задачи 't2' и 't3' будут запускаться одновременно после завершения задачи 't1'. Затем задача 't4' запустится после завершения обеих задач 't2' и 't3'.
|
||||
|
||||
```python
|
||||
t1 >> [t2, t3] # t2 и t3 запускаются одновременно после t1
|
||||
[t2, t3] >> t4 # t4 запускается после завершения t2 и t3
|
||||
```
|
||||
|
||||
Для сложных пайплайнов можно создавать цепочки любой сложности:
|
||||
```python
|
||||
t1 >> [t2, t3] >> t4 >> [t5, t6, t7] >> t8
|
||||
```
|
||||
|
||||
## Полный пример DAG-файла
|
||||
|
||||
В этом полном примере мы создаем DAG с более подробной настройкой параметров. Обратите внимание на дополнительные параметры, такие как количество повторных попыток ('retries'), задержка между попытками ('retry_delay'), а также настройки уведомлений по электронной почте.
|
||||
|
||||
Вот как выглядит готовый DAG-файл в итоге:
|
||||
|
||||
```python
|
||||
from datetime import datetime, timedelta
|
||||
from airflow import DAG
|
||||
from airflow.operators.bash import BashOperator
|
||||
from airflow.utils.dates import days_ago
|
||||
|
||||
|
||||
default_args = {
|
||||
'owner': 'data_team',
|
||||
'depends_on_past': False,
|
||||
'start_date': days_ago(1),
|
||||
'email': ['data@example.com'],
|
||||
'email_on_failure': False,
|
||||
'email_on_retry': False,
|
||||
'retries': 2,
|
||||
'retry_delay': timedelta(minutes=3),
|
||||
}
|
||||
|
||||
dag = DAG(
|
||||
"sample_dag",
|
||||
default_args=default_args,
|
||||
schedule_interval=None,
|
||||
)
|
||||
|
||||
t1 = BashOperator(
|
||||
task_id='show_time',
|
||||
bash_command='date',
|
||||
dag=dag
|
||||
)
|
||||
|
||||
t2 = BashOperator(
|
||||
task_id='process_data',
|
||||
bash_command='sleep 3',
|
||||
retries=2,
|
||||
dag=dag
|
||||
)
|
||||
|
||||
t1 >> t2
|
||||
```
|
||||
|
||||
**Важное замечание**: код DAG-файла выполняется Airflow при каждом сканировании директории DAG. Поэтому в нем не следует размещать тяжелые вычисления, чтение больших файлов или запросы к базам данных. DAG-файл должен быть максимально легковесным. Сложные операции следует выносить в отдельные функции и вызывать их через соответствующие операторы.
|
||||
|
||||
💡 **Правило**: DAG-файл не должен содержать ресурсоемких операций, загрузки файлов или сложных вычислений.
|
||||
|
||||
# Создаем свой первый рабочий пайплайн
|
||||
|
||||
Теперь давайте создадим полноценный пайплайн для работы с реальными данными. В качестве примера мы будем использовать датасет с информацией о клиентах, который часто применяется в задачах анализа данных.
|
||||
|
||||
Наш пайплайн будет состоять из трех этапов:
|
||||
1. Загрузка исходного датасета из интернета
|
||||
2. Создание агрегированной таблицы по регионам и категориям
|
||||
3. Сохранение результата в базу данных PostgreSQL
|
||||
|
||||
Вот полный код нашего DAG:
|
||||
|
||||
```python
|
||||
import os
|
||||
import datetime as dt
|
||||
import pandas as pd
|
||||
from airflow.models import DAG
|
||||
from airflow.operators.python import PythonOperator
|
||||
from airflow.operators.bash import BashOperator
|
||||
|
||||
from sqlalchemy import create_engine
|
||||
|
||||
|
||||
# Базовые параметры DAG
|
||||
args = {
|
||||
'owner': 'airflow',
|
||||
'start_date': dt.datetime(2020, 12, 23),
|
||||
'retries': 1,
|
||||
'retry_delay': dt.timedelta(minutes=1),
|
||||
}
|
||||
|
||||
def download_titanic_dataset():
|
||||
url = 'https://web.stanford.edu/class/archive/cs/cs109/cs109.1166/stuff/titanic.csv'
|
||||
df = pd.read_csv(url)
|
||||
engine = create_engine('postgresql+psycopg2://jovyan:jovyan@localhost:5432/de')
|
||||
df.to_sql('titanic', engine, index=False, if_exists='replace', schema='public')
|
||||
|
||||
|
||||
def pivot_dataset():
|
||||
engine = create_engine('postgresql+psycopg2://jovyan:jovyan@localhost:5432/de')
|
||||
titanic_df = pd.read_sql('select * from public.titanic', con=engine)
|
||||
|
||||
df = titanic_df.pivot_table(
|
||||
index=['Sex'],
|
||||
columns=['Pclass'],
|
||||
values='Name',
|
||||
aggfunc='count'
|
||||
).reset_index()
|
||||
|
||||
df.to_sql('titanic_pivot', engine, index=False, if_exists='replace', schema='public' )
|
||||
|
||||
dag = DAG(
|
||||
dag_id='titanic_pivot',
|
||||
schedule_interval=None,
|
||||
default_args=args,
|
||||
)
|
||||
|
||||
# Начальная задача для логирования
|
||||
start = BashOperator(
|
||||
task_id='start',
|
||||
bash_command='echo "Начинаем выполнение пайплайна! "',
|
||||
dag=dag,
|
||||
)
|
||||
|
||||
# Загрузка исходного датасета
|
||||
create_titanic_dataset = PythonOperator(
|
||||
task_id='download_titanic_dataset',
|
||||
python_callable=download_titanic_dataset,
|
||||
dag=dag,
|
||||
)
|
||||
|
||||
# Преобразование и сохранение сводной таблицы
|
||||
pivot_titanic_dataset = PythonOperator(
|
||||
task_id='pivot_dataset',
|
||||
python_callable=pivot_dataset,
|
||||
dag=dag,
|
||||
)
|
||||
|
||||
# Последовательность выполнения
|
||||
start >> create_titanic_dataset >> pivot_titanic_dataset
|
||||
```
|
||||
|
||||
Этот код использует `PythonOperator` для выполнения функций работы с данными, что является стандартной практикой для задач обработки данных. В примере функция `load_customer_dataset` загружает данные из внешнего источника, а `aggregate_customer_dataset` создает агрегированную таблицу.
|
||||
|
||||
## Загрузка DAG в учебную среду
|
||||
|
||||
Для тестирования нашего DAG в учебной среде выполните следующие шаги:
|
||||
|
||||
1. Сохраните код в файл с расширением `.py` (например, `customer_analysis_dag.py`)
|
||||
|
||||
Для тестирования нашего DAG в учебной среде выполните следующие шаги:
|
||||
|
||||
1. Сохраните код в файл с расширением `.py` (например, `customer_analysis_dag.py`)
|
||||
|
||||
2. Найдите запущенный контейнер с учебной средой Airflow:
|
||||
```bash
|
||||
docker ps
|
||||
```
|
||||
|
||||
3. Скопируйте файл в контейнер:
|
||||
```bash
|
||||
docker cp /путь/к/файлу/customer_analysis_dag.py [ID_КОНТЕЙНЕРА]:/lessons/dags/customer_analysis_dag.py
|
||||
```
|
||||
|
||||
Например:
|
||||
```bash
|
||||
docker cp ~/Desktop/customer_analysis_dag.py 4dc4fbfedcf4:/lessons/dags/customer_analysis_dag.py
|
||||
```
|
||||
|
||||
4. Подождите 30-60 секунд — Airflow автоматически обнаружит новый файл
|
||||
|
||||
5. Найдите ваш DAG в интерфейсе Airflow через строку поиска и запустите его
|
||||
|
||||

|
||||
|
||||
В этом уроке вы изучили основную структуру DAG-файлов, создали свой первый рабочий пайплайн, научились загружать его в учебную среду и запускать для получения результатов.
|
||||
@@ -0,0 +1,78 @@
|
||||
# Статусы задач в Airflow
|
||||
|
||||
# Статусы задач в Airflow — что означают цвета и как ими пользоваться
|
||||
|
||||
## Почему статусы задач так важны?
|
||||
|
||||
Когда вы только начинаете работать с Airflow, интерфейс может показаться сложным. Особенно если вы видите задачи разных цветов и не понимаете, что это значит. На самом деле, статусы задач — это ваш главный помощник в отладке и понимании того, что происходит с вашим пайплайном.
|
||||
|
||||
## Основные статусы, которые вы увидите каждый день
|
||||
|
||||
В интерфейсе Airflow каждая задача отображается определенным цветом. Вот что означают самые важные статусы:
|
||||
|
||||

|
||||
|
||||
### Простое объяснение всех статусов
|
||||
|
||||
Давайте разберем каждый статус простым языком:
|
||||
|
||||
**🟢 Успешно (success)** — ваша задача выполнилась без ошибок. Это то, к чему мы стремимся!
|
||||
|
||||
**🔴 Ошибка (failed)** — что-то пошло не так. Задача упала, и вам нужно разбираться в коде.
|
||||
|
||||
**🟡 В очереди (queued)** — задача ждет своей очереди на выполнение. Это нормально, особенно если у вас много задач или мало ресурсов.
|
||||
|
||||
**🔵 Выполняется (running)** — задача сейчас активно работает. Просто подождите немного.
|
||||
|
||||
**⚪ Нет статуса (none/no status)** — задача еще не готова к запуску, потому что не выполнены её зависимости.
|
||||
|
||||
**🟠 Запланирована (scheduled)** — все готово к запуску, Airflow вот-вот начнет выполнение.
|
||||
|
||||
**🟣 Пропущена (skipped)** — задача была намеренно пропущена (часто в ветвящихся пайплайнах).
|
||||
|
||||
### Специальные статусы (встречаются реже)
|
||||
|
||||
- **Ошибка в зависимости (upstream_failed)** — предыдущая задача упала, поэтому текущая даже не запускалась
|
||||
- **Готова к повтору (up_for_retry)** — задача упала, но Airflow попробует запустить её снова (если настроены повторные попытки)
|
||||
- **Завершена (shutdown)** — задачу принудительно остановили во время выполнения
|
||||
- **Отложена (deferred)** — задача приостановлена и ждет внешнего события
|
||||
- **Наблюдение (sensing)** — специальный статус для сенсоров, которые ждут определенных условий
|
||||
|
||||
## Как задача проходит свой путь: пошагово
|
||||
|
||||
Представьте, что у вас есть простая задача. Вот как она проходит свой жизненный цикл:
|
||||
|
||||
1. **Создание** → Статус: "Нет статуса"
|
||||
Airflow создает задачу, но еще не может её запустить
|
||||
|
||||
2. **Готовность** → Статус: "Запланирована"
|
||||
Все зависимости выполнены, задача готова к работе
|
||||
|
||||
3. **Ожидание** → Статус: "В очереди"
|
||||
Задача ждет свободного рабочего места
|
||||
|
||||
4. **Работа** → Статус: "Выполняется"
|
||||
Задача активно выполняется
|
||||
|
||||
5. **Завершение** → Статус: "Успешно"
|
||||
Всё прошло отлично!
|
||||
|
||||

|
||||
|
||||
## Что делать, если задача упала?
|
||||
|
||||
Если вы видите красный статус (failed), не паникуйте! Это нормальная часть работы с данными. Вот что делать:
|
||||
|
||||
1. **Нажмите на задачу** в интерфейсе Airflow
|
||||
2. **Посмотрите логи** — там будет точная причина ошибки
|
||||
3. **Исправьте код** или настройки
|
||||
4. **Перезапустите задачу** кнопкой "Clear"
|
||||
|
||||
## Советы для начинающих
|
||||
|
||||
- **Не бойтесь статусов** — они ваш друг, а не враг
|
||||
- **Самые важные статусы** для начала: success, failed, queued, running
|
||||
- **Остальные статусы** вы будете изучать по мере необходимости
|
||||
- **Цвета в интерфейсе** — это быстрый способ понять состояние вашего пайплайна
|
||||
|
||||
Помните: понимание статусов задач — это как научиться читать дорожные знаки. Сначала кажется много информации, но со временем это становится второй натурой!
|
||||
@@ -0,0 +1,142 @@
|
||||
# Гибкие шаблоны и настройки в Airflow
|
||||
В этом материале вы познакомитесь с мощными инструментами Airflow для создания гибких и переиспользуемых пайплайнов: динамическими шаблонами, безопасными переменными и централизованными подключениями к внешним системам.
|
||||
|
||||
# Динамические шаблоны Airflow (на основе Jinja)
|
||||
|
||||
В процессах обработки данных часто возникает необходимость использовать переменные значения, такие как дата выполнения задачи, для фильтрации исходных данных. В сложных сценариях может потребоваться информация о предыдущих успешных запусках или метаданные текущего выполнения. Airflow предоставляет встроенный механизм динамических шаблонов, основанный на языке Jinja, который значительно упрощает работу с такими переменными.
|
||||
|
||||
Синтаксис шаблонов интуитивно понятен — переменная заключается в двойные фигурные скобки с пробелами по краям. Для эффективного использования достаточно знать, какое значение возвращает конкретный шаблон, и разместить его в нужном месте кода.
|
||||
|
||||
Подробный справочник доступных шаблонов можно найти в [официальной документации Airflow](https://airflow.apache.org/docs/apache-airflow/stable/templates-ref.html#variables).
|
||||
|
||||
Вот практические примеры использования:
|
||||
|
||||
**Передача даты выполнения в задачу:**
|
||||
|
||||
В этом примере создается DAG с идентификатором "template_example", который запускается каждые 15 минут. Задача "display_date" использует шаблон {{ ds }} для вывода текущей даты выполнения в формате YYYY-MM-DD.
|
||||
```python
|
||||
...
|
||||
|
||||
dag = DAG(
|
||||
dag_id="dynamic_templates_example",
|
||||
schedule_interval="*/10 * * * *",
|
||||
default_args=default_args
|
||||
)
|
||||
|
||||
t1 = BashOperator(task_id="show_date", bash_command="echo {{ ds }}")
|
||||
|
||||
t1
|
||||
```
|
||||
|
||||
**Использование шаблонов в идентификаторах задач для лучшей отслеживаемости:**
|
||||
|
||||
В этом примере создается DAG с идентификатором "template_tracking_example", который запускается каждые 20 минут. Вторая задача использует шаблон {{ ds }} в своем идентификаторе, что позволяет легко идентифицировать задачу по дате выполнения.
|
||||
```python
|
||||
...
|
||||
|
||||
dag = DAG(
|
||||
dag_id="dynamic_templates_example",
|
||||
schedule_interval="*/10 * * * *",
|
||||
default_args=default_args
|
||||
)
|
||||
|
||||
t1 = BashOperator(task_id="show_date", bash_command="echo {{ ds }}")
|
||||
t2 = BashOperator(task_id=f"process_for_{{ ds }}", bash_command="echo {{ ds }}")
|
||||
|
||||
t1 >> t2
|
||||
```
|
||||
|
||||
# Безопасные переменные (Variables)
|
||||
|
||||
Переменные Airflow представляют собой пары "ключ-значение", хранящиеся в метадатабазе системы. Они идеально подходят для хранения конфигурационных параметров, таких как пути к скриптам, имена таблиц или другие настройки, которые должны быть доступны в разных DAG.
|
||||
|
||||
Управление переменными осуществляется через веб-интерфейс Airflow (Admin → Variables), где можно:
|
||||
- Создавать и редактировать пары ключ-значение вручную
|
||||
- Импортировать настройки из JSON-файлов
|
||||
- Использовать командную строку Airflow
|
||||
|
||||

|
||||
|
||||
Для защиты конфиденциальной информации Airflow автоматически маскирует значения переменных, в названии которых содержится слово "secret".
|
||||
|
||||
**Пример использования переменной в коде DAG:**
|
||||
|
||||
В этом примере создается DAG с идентификатором "variable_example", который использует переменную 'data_storage_path', предварительно сохраненную в Airflow. Значение переменной извлекается с помощью Variable.get() и используется в команде bash для указания пути к данным.
|
||||
```python
|
||||
from airflow import DAG
|
||||
from airflow.operators.bash import BashOperator
|
||||
from airflow.operators.dummy import DummyOperator
|
||||
from airflow.models import Variable
|
||||
from datetime import datetime
|
||||
|
||||
default_args = {
|
||||
'owner': 'data_team',
|
||||
'start_date': datetime(2023, 1, 1),
|
||||
'retries': 1,
|
||||
}
|
||||
|
||||
dag = DAG(
|
||||
dag_id="variable_example",
|
||||
schedule_interval=None,
|
||||
default_args=default_args
|
||||
)
|
||||
|
||||
# Получение значения переменной
|
||||
data_path = Variable.get('data_storage_path')
|
||||
|
||||
task = BashOperator(
|
||||
task_id='process_data',
|
||||
bash_command=f'echo "Processing data from {data_path}"',
|
||||
dag=dag
|
||||
)
|
||||
```
|
||||
|
||||
Переменные делают код DAG более читаемым и модульным, позволяя легко адаптировать один и тот же пайплайн для разных окружений или сценариев использования.
|
||||
|
||||
# Централизованные подключения (Connections)
|
||||
|
||||
Подключения в Airflow — это безопасный способ хранения учетных данных и параметров для взаимодействия с внешними системами. Каждое подключение имеет уникальный идентификатор (conn_id) и содержит необходимые параметры: хост, порт, логин, пароль и другие специфичные настройки.
|
||||
|
||||
Подключения поддерживают широкий спектр систем:
|
||||
- Базы данных (PostgreSQL, MySQL, Oracle и др.)
|
||||
- Облачные хранилища (AWS S3, Google Cloud Storage)
|
||||
- Системы уведомлений (Email, Telegram)
|
||||
- И многие другие через Airflow Providers
|
||||
|
||||
При использовании операторов, взаимодействующих с внешними системами, Airflow автоматически ищет соответствующее подключение с суффиксом `_default`. Однако можно явно указать альтернативное подключение:
|
||||
|
||||
В приведенном примере используется PostgresOperator для создания таблицы в базе данных. Вместо использования подключения по умолчанию, явно указывается подключение с идентификатором 'my_postgres_conn'.
|
||||
|
||||
```python
|
||||
from airflow.operators.postgres_operator import PostgresOperator
|
||||
|
||||
create_table = PostgresOperator(
|
||||
task_id='create_user_table',
|
||||
sql='''
|
||||
CREATE TABLE users(
|
||||
user_id integer NOT NULL,
|
||||
created_at TIMESTAMP NOT NULL
|
||||
);''',
|
||||
postgres_conn_id='my_postgres_conn'
|
||||
)
|
||||
```
|
||||
|
||||
Управление подключениями доступно через интерфейс Airflow (Admin → Connections). Если требуемый тип подключения отсутствует, его можно добавить установкой соответствующего Airflow Provider из [официального репозитория](https://airflow.apache.org/docs/#providers-packages-docs-apache-airflow-providers-index-html).
|
||||
|
||||

|
||||
|
||||
Для успешной работы с внешними системами сначала необходимо создать соответствующее подключение, а затем использовать его идентификатор в операторах вашего DAG.
|
||||
|
||||
Более подробную информацию о настройке подключений можно найти в [документации Airflow](https://airflow.apache.org/docs/apache-airflow/stable/howto/connection.html).
|
||||
|
||||
# Проверочный список для качественного DAG
|
||||
|
||||
После создания DAG задайте себе следующие вопросы для обеспечения его качества и безопасности:
|
||||
|
||||
- Сможет ли коллега понять и поддерживать этот DAG в моё отсутствие?
|
||||
- Содержит ли код чувствительную информацию (логины, пароли, API-ключи)?
|
||||
- Какие параметры можно вынести в переменные для лучшей гибкости?
|
||||
- Требуется ли маскировка конфиденциальных значений?
|
||||
- Используются ли в логике даты или временные метки, которые можно заменить на шаблоны?
|
||||
|
||||
Ответы на эти вопросы помогут вам создавать надежные, безопасные и легко поддерживаемые пайплайны в Airflow.
|
||||
@@ -0,0 +1,96 @@
|
||||
# Управление временем
|
||||
|
||||
В этом практическом руководстве вы освоите ключевые аспекты работы со временем в Airflow. Вы узнаете, как правильно настраивать расписания для ваших DAG, понимать внутреннюю логику временных меток и эффективно управлять историческими данными. Поскольку бизнес-процессы часто напрямую зависят от временных параметров запуска, понимание этих механизмов критически важно для любого дата-инженера.
|
||||
|
||||
# Настройка автоматических запусков
|
||||
|
||||
Когда вы создаете свой первый пайплайн, вы уже сталкивались с возможностью запускать DAG по расписанию через параметр `schedule_interval`. По умолчанию этот параметр равен `None`, что означает ручной запуск без автоматического расписания.
|
||||
|
||||
В приведенном примере создается DAG с идентификатором "daily_data_processing", который будет запускаться ежедневно. DAG использует дату начала 1 января 2020 года и параметр catchup=False, что означает, что пропущенные запуски обрабатываться не будут. Владелец DAG - команда data_team, и для задач в DAG установлена одна попытка повторного запуска при ошибках.
|
||||
|
||||
```python
|
||||
dag = DAG(
|
||||
dag_id="daily_data_processing",
|
||||
schedule_interval="@daily",
|
||||
start_date=dt.datetime(2020, 1, 1),
|
||||
catchup=False,
|
||||
default_args={
|
||||
'owner': 'data_team',
|
||||
'retries': 1,
|
||||
}
|
||||
)
|
||||
```
|
||||
|
||||
Важно понимать, что Airflow требует указания даты начала работы (`start_date`) для любого DAG с расписанием. Система использует эту дату как отправную точку и планирует первый запуск, добавляя к ней интервал из `schedule_interval`.
|
||||
|
||||
## Готовые шаблоны расписания
|
||||
|
||||
Airflow предоставляет удобные встроенные шаблоны для самых распространенных сценариев:
|
||||
|
||||
- `@once` — однократный запуск
|
||||
- `@hourly` — каждый час
|
||||
- `@daily` — ежедневно
|
||||
- `@weekly` — еженедельно
|
||||
- `@monthly` — ежемесячно
|
||||
- `@yearly` — ежегодно
|
||||
|
||||
Для более сложных сценариев можно использовать стандартный cron-формат, например: `"0 12 * * 1-5"` для запуска в 12:00 по будням.
|
||||
|
||||
# Работа с историческими данными
|
||||
|
||||
## Механизм "догонки" (catchup)
|
||||
|
||||
Когда вы создаете DAG с исторической датой начала, Airflow предлагает мощный механизм автоматического пересчета пропущенных периодов через параметр `catchup`.
|
||||
|
||||
В этом примере создается DAG с идентификатором "historical_data_processing", который запускается ежедневно и имеет дату начала 1 января 2021 года. Параметр catchup=True означает, что Airflow будет автоматически запускать DAG для всех пропущенных дней с указанной даты начала до текущего момента. Владелец DAG - команда analytics_team, и для задач установлено две попытки повторного запуска при ошибках.
|
||||
|
||||
Когда вы создаете DAG с исторической датой начала, Airflow предлагает мощный механизм автоматического пересчета пропущенных периодов через параметр `catchup`.
|
||||
|
||||
```python
|
||||
dag = DAG(
|
||||
dag_id="historical_data_processing",
|
||||
schedule_interval="@daily",
|
||||
start_date=dt.datetime(2021, 1, 1),
|
||||
catchup=True,
|
||||
default_args={
|
||||
'owner': 'analytics_team',
|
||||
'retries': 2,
|
||||
}
|
||||
)
|
||||
```
|
||||
|
||||
При `catchup=True` система автоматически выполнит все пропущенные запуски от указанной даты начала до текущего момента. Это особенно полезно при первом запуске DAG для обработки накопившихся исторических данных.
|
||||
|
||||
Если ваш бизнес-сценарий не требует пересчета истории или вы хотите начать обработку только с текущего периода, установите `catchup=False`. В этом случае Airflow будет планировать только ближайшие запуски согласно расписанию.
|
||||
|
||||
## Ручная перезаливка данных (backfill)
|
||||
|
||||
Иногда возникает необходимость пересчитать данные за конкретный период времени. Для этого Airflow предоставляет команду `backfill`:
|
||||
|
||||
```bash
|
||||
airflow dags backfill \
|
||||
--start-date 2022-01-01 \
|
||||
--end-date 2022-03-01 \
|
||||
the_main_dag
|
||||
```
|
||||
|
||||
Эта команда запустит DAG `the_main_dag` для каждого интервала между указанными датами, позволяя гибко управлять перерасчетом исторических данных.
|
||||
|
||||
# Временные зоны и локализация
|
||||
|
||||
Важный момент: **все вычисления в Airflow по умолчанию выполняются в UTC** (на 3 часа меньше московского времени). Это стандартная практика для распределенных систем, но требует особого внимания при работе с локальными временными метками.
|
||||
|
||||
Вы можете:
|
||||
- Настроить глобальную временную зону в конфигурационном файле Airflow
|
||||
- Указать временную зону явно при инициализации DAG через параметр `tz`
|
||||
|
||||
Однако рекомендуется придерживаться UTC во всех расчетах и преобразовывать временные метки только при выводе результатов для конечных пользователей. Это минимизирует ошибки и упрощает отладку.
|
||||
|
||||
# Практические рекомендации
|
||||
|
||||
1. **Всегда тестируйте расписание** на небольшом временном интервале перед запуском в продакшен
|
||||
2. **Используйте `catchup=False`** для DAG, которые не требуют исторических пересчетов
|
||||
3. **Планируйте запуски с учетом UTC**, особенно если ваша команда работает в разных часовых поясах
|
||||
4. **Документируйте временные зависимости** в коде DAG для других разработчиков
|
||||
|
||||
Понимание временных механизмов Airflow — ключ к созданию надежных и предсказуемых пайплайнов. Правильная настройка расписаний и управление историческими данными позволяют автоматизировать сложные бизнес-процессы без ручного вмешательства.
|
||||
@@ -0,0 +1,429 @@
|
||||
# Продвинутые возможности Airflow для начинающих специалистов
|
||||
|
||||
В этом материале мы познакомимся с расширенными функциями Airflow, которые помогут вам решать нетривиальные задачи при работе со сложными пайплайнами.
|
||||
|
||||
Начнем с улучшения нашего базового пайплайна — проведем рефакторинг кода для лучшей читаемости и поддержки.
|
||||
|
||||
## Управление ресурсами с помощью пулов задач
|
||||
|
||||
В системах с высокой нагрузкой, где одновременно запускается множество задач и DAG-ов, может возникнуть чрезмерная нагрузка на исполнителей и серверную часть. Это может привести к ошибкам выполнения и даже к отказу системы, если не установить соответствующие ограничения.
|
||||
|
||||
В Airflow для решения этой проблемы существует механизм управления ресурсами — **пулы задач** (pools). По умолчанию в Airflow настроен один пул задач — `default_pool` с 128 слотами, что означает возможность параллельного выполнения 128 задач одновременно. Пул `default_pool` нельзя удалить, но можно изменить его размер — увеличить или уменьшить количество слотов.
|
||||
|
||||
Когда планировщик обнаруживает, что наступило время выполнения DAG, он запускает задачу согласно заданной последовательности. При этом задача занимает один слот в пуле и освобождает его после завершения.
|
||||
|
||||
Создать новый пул задач и установить его размер можно через веб-интерфейс Airflow:
|
||||
|
||||

|
||||
|
||||
*Создание нового пула задач в Airflow*
|
||||
|
||||
Зачем нужны пулы задач? Они помогают:
|
||||
- Организовать запуск процессов в системе
|
||||
- Предотвратить перегрузку системы при выполнении большого количества ресурсоемких задач
|
||||
|
||||
Вы можете задать "вес" задачи через параметр `pool_slots`, чтобы оптимизировать распределение нагрузки. Если общее количество задач превышает доступные слоты, планировщик поставит задачу в очередь и запустит ее, как только появятся свободные ресурсы.
|
||||
|
||||
Пример настройки веса задач:
|
||||
|
||||
В приведенном примере показано, как можно настроить использование пула задач с различным весом. Задача 'data_backup_task' использует 3 слота в пуле 'data_processing_pool', что делает ее более ресурсоемкой по сравнению с задачами 'file_check_task' и 'cleanup_task', которые используют по 1 слоту. Это позволяет контролировать распределение ресурсов между различными задачами.
|
||||
|
||||
```python
|
||||
BashOperator(
|
||||
task_id="data_backup_task",
|
||||
bash_command="bash backup_script.sh",
|
||||
pool_slots=3,
|
||||
pool="data_processing_pool",
|
||||
)
|
||||
|
||||
BashOperator(
|
||||
task_id="file_check_task",
|
||||
bash_command="bash validate_files.sh",
|
||||
pool_slots=1,
|
||||
pool="data_processing_pool",
|
||||
)
|
||||
|
||||
BashOperator(
|
||||
task_id="cleanup_task",
|
||||
bash_command="bash cleanup_files.sh",
|
||||
pool_slots=1,
|
||||
pool="data_processing_pool",
|
||||
)
|
||||
```
|
||||
|
||||
Более подробную информацию о механизме пулов можно найти в [официальной документации](https://airflow.apache.org/docs/apache-airflow/stable/concepts/pools.html#pools).
|
||||
|
||||
## Обмен данными между задачами: XCom и контекст выполнения
|
||||
|
||||
**XCom** (от англ. cross-communications — «межзадачная коммуникация») — это механизм обмена сообщениями между задачами внутри одного DAG. Архитектурно каждая задача изолирована от других и работает в собственном контексте.
|
||||
|
||||
💡 **Контекст выполнения** — это набор параметров, передаваемых при запуске задачи, а также метаданные, генерируемые во время выполнения: время старта, время завершения, имя DAG и другие. Контекст можно найти следующим образом:
|
||||
|
||||
{:height 437, :width 778}
|
||||
|
||||
В XCom можно передавать сериализованные объекты. Значения XCom хранятся в базе данных Airflow и доступны через интерфейс:
|
||||
|
||||

|
||||
|
||||
*XCom в веб-интерфейсе Airflow*
|
||||
|
||||
**Важное правило**: XCom предназначен для обмена небольшими сообщениями. Данные проходят сериализацию/десериализацию при чтении и записи в таблицу.
|
||||
|
||||
💡 **Сериализация** — процесс преобразования структуры данных в последовательность байтов. **Десериализация** — восстановление структуры данных из байтовой последовательности.
|
||||
|
||||
Для передачи больших объемов данных используйте внешние средства: файловую систему, базы данных (чаще всего PostgreSQL) или очереди сообщений (например, Kafka).
|
||||
|
||||
Механизм XCom похож на работу функций в Python. Многие операторы по умолчанию возвращают результат выполнения задачи. За это отвечает параметр `do_xcom_push`, который по умолчанию равен `True`.
|
||||
|
||||
Чтобы прочитать сообщения из XCom, используйте метод `xcom_pull` в контексте задачи:
|
||||
|
||||
В этом примере мы получаем результат выполнения задачи с идентификатором 'data_processing_task' с помощью метода xcom_pull. Это позволяет передавать небольшие объемы данных между задачами в рамках одного DAG.
|
||||
|
||||
```python
|
||||
result = task_instance.xcom_pull(task_ids='data_processing_task')
|
||||
```
|
||||
|
||||
Также можно обращаться к сообщениям через Jinja-шаблоны:
|
||||
|
||||
```
|
||||
SELECT * FROM {{ task_instance.xcom_pull(task_ids='foo', key='table_name') }}
|
||||
```
|
||||
|
||||
Если оператор возвращает значение и параметр `do_xcom_push` установлен в `True` (по умолчанию), это значение автоматически записывается в XCom.
|
||||
|
||||
Пример явного запрета записи в XCom:
|
||||
|
||||
В этом примере задача 'show_directory_contents' создается с параметром do_xcom_push=False, что означает, что результат выполнения этой задачи не будет автоматически сохранен в XCom. Это полезно, когда вы не хотите, чтобы задача передавала какие-либо данные другим задачам через XCom.
|
||||
|
||||
```python
|
||||
list_files = BashOperator(
|
||||
task_id='show_directory_contents',
|
||||
bash_command='ls -la',
|
||||
do_xcom_push=False
|
||||
)
|
||||
|
||||
def calculate_sum():
|
||||
return 2 + 3
|
||||
```
|
||||
|
||||
Пример автоматической записи в XCom (параметр `do_xcom_push` по умолчанию `True`):
|
||||
|
||||
В этом примере задача 'calculate_sum' автоматически записывает результат выполнения функции calculate_sum в XCom, так как параметр do_xcom_push по умолчанию установлен в True. Это позволяет использовать результат этой задачи в других задачах DAG через XCom.
|
||||
|
||||
```python
|
||||
sum_result = PythonOperator(
|
||||
task_id='calculate_sum',
|
||||
python_callable=calculate_sum,
|
||||
)
|
||||
```
|
||||
|
||||
XCom похожи на переменные (variables) в Airflow, но предназначены именно для взаимодействия между задачами в рамках одного DAG, а не для глобальных настроек.
|
||||
|
||||
XCom упрощает взаимодействие между задачами и применяется в различных сценариях.
|
||||
|
||||
## Группировка задач с помощью TaskGroup
|
||||
|
||||
Для удобства визуализации задачи можно группировать в веб-интерфейсе Airflow (начиная с версии 2.0). Повторяющиеся или логически связанные задачи можно объединить в группы:
|
||||
|
||||
В этом примере задачи сгруппированы в три логические группы: 'data_extraction', 'data_transformation' и 'data_loading'. Каждая группа содержит несколько задач, которые выполняются последовательно внутри группы. Затем группы связаны между собой, чтобы показать общий порядок выполнения этапов обработки данных.
|
||||
|
||||
```python
|
||||
with TaskGroup("data_extraction") as extraction_group:
|
||||
extract_1 = DummyOperator(task_id="extract_source_1")
|
||||
extract_2 = DummyOperator(task_id="extract_source_2")
|
||||
extract_3 = DummyOperator(task_id="extract_source_3")
|
||||
extract_1 >> extract_2 >> extract_3
|
||||
|
||||
with TaskGroup("data_transformation") as transformation_group:
|
||||
transform_1 = DummyOperator(task_id="transform_step_1")
|
||||
transform_2 = DummyOperator(task_id="transform_step_2")
|
||||
transform_1 >> transform_2
|
||||
|
||||
with TaskGroup("data_loading") as loading_group:
|
||||
load_1 = DummyOperator(task_id="load_to_target_1")
|
||||
load_2 = DummyOperator(task_id="load_to_target_2")
|
||||
load_1 >> load_2
|
||||
|
||||
[extraction_group, transformation_group] >> loading_group
|
||||
```
|
||||
|
||||
Обратите внимание: при использовании TaskGroup последовательность задач указывается внутри группы после объявления всех задач, а в конце DAG описывается последовательность выполнения самих групп.
|
||||
|
||||
Визуально в интерфейсе Airflow это выглядит так:
|
||||
|
||||

|
||||
|
||||
*Группировка задач с помощью TaskGroup*
|
||||
|
||||
TaskGroup добавляет интерактивность в веб-интерфейс — группы задач можно сворачивать и разворачивать, что значительно улучшает восприятие DAG с большим количеством задач и связей:
|
||||
|
||||

|
||||
|
||||
*Группу задач можно раскрыть для детального просмотра*
|
||||
|
||||
TaskGroup — это удобный способ логической группировки задач, который помогает упростить код и представить сложные пайплайны более компактно.
|
||||
|
||||
## Система оповещений (алертинг)
|
||||
|
||||
Алертинг — один из ключевых компонентов системы оркестрации, так как важно своевременно получать уведомления об ошибках для их оперативного анализа и решения.
|
||||
|
||||
По умолчанию в Airflow настроена отправка уведомлений на электронную почту. При создании DAG указываются email-адреса, на которые будут отправляться сообщения. С помощью параметров можно настроить различные сценарии оповещений.
|
||||
|
||||
Давайте модифицируем наш первый DAG так, чтобы получать уведомления на почту при возникновении ошибок. При этом настроим перезапуск задач в случае неудачи (например, 2 попытки), но без уведомлений о самих перезапусках:
|
||||
|
||||
В приведенном примере создан DAG с идентификатором 'customer_analysis_pipeline', который настроен на отправку уведомлений по электронной почте только при ошибках (email_on_failure=True), но не при повторных попытках (email_on_retry=False). Также установлено 2 попытки повторного запуска задач при ошибках с задержкой 2 минуты между попытками.
|
||||
|
||||
```python
|
||||
import os
|
||||
import datetime as dt
|
||||
import pandas as pd
|
||||
from airflow.models import DAG
|
||||
from airflow.operators.python import PythonOperator
|
||||
from airflow.operators.bash import BashOperator
|
||||
from sqlalchemy import create_engine
|
||||
|
||||
# основные параметры DAG
|
||||
args = {
|
||||
'owner': 'data_engineering_team',
|
||||
'start_date': dt.datetime(2021, 6, 15),
|
||||
'retries': 2,
|
||||
'retry_delay': dt.timedelta(minutes=2),
|
||||
'email': ["data-team@example.com"],
|
||||
'email_on_failure': True,
|
||||
'email_on_retry': False,
|
||||
}
|
||||
|
||||
dag = DAG(
|
||||
dag_id='customer_analysis_pipeline',
|
||||
schedule_interval=None,
|
||||
default_args=args,
|
||||
)
|
||||
```
|
||||
|
||||
Теперь вы будете получать email-уведомления при ошибках выполнения.
|
||||
|
||||
Пример функции для отправки уведомлений с использованием параметров из контекста:
|
||||
|
||||
```python
|
||||
from datetime import datetime, timedelta, timezone
|
||||
import dateutil
|
||||
from airflow.utils.email import send_email_smtp
|
||||
|
||||
MAIL_LIST = [
|
||||
"email_1@gmail.ru",
|
||||
"email_2@gmail.ru"
|
||||
]
|
||||
|
||||
def notify_email(calculation_dt: str, dagrun_begin_time, **context):
|
||||
# dag_run date is utc timezone, so add `timezone.utc` to calculation_end to combat 3 hour diff.
|
||||
calculation_start = dateutil.parser.isoparse(dagrun_begin_time.replace('Z', '+00:00'))
|
||||
calculation_end = datetime.now(timezone.utc)
|
||||
duration = str(timedelta(seconds=(calculation_end - calculation_start).seconds))
|
||||
|
||||
title = f"Ежедневные расчёты завершились успешно (dag: {context['task_instance'].dag_id})."
|
||||
|
||||
body = f"""
|
||||
Привет, <br>
|
||||
Я закончил работу над расчётом за "{calculation_dt}". Это заняло {duration} часов/минут/секунд.<br>
|
||||
Логи и запуски тоже можно посмотреть <a href="http://airflow-monitoring.example.com/tree?dag_id=daily_guests_features">тут</a>.
|
||||
<br>
|
||||
|
||||
<br>
|
||||
<br>
|
||||
Навеки твой,<br>
|
||||
Airflow бот <br>
|
||||
"""
|
||||
|
||||
send_email_smtp(";".join(MAIL_LIST), title, body)
|
||||
```
|
||||
|
||||
Такую функцию можно разместить в отдельном файле (например, `utils.py`), импортировать как модуль в нужных DAG и вызывать отдельной задачей:
|
||||
|
||||
```python
|
||||
from airflow.operators.python import PythonOperator
|
||||
from utils import notify_email
|
||||
|
||||
LOCAL_CALCULATION_DT = '{{ dag.timezone.convert(execution_date).strftime("%Y-%m-%d") }}'
|
||||
DAG_RUN_BEGIN_TIME = "{{ dag_run.start_date }}"
|
||||
|
||||
email_notification_task = PythonOperator(
|
||||
task_id="send_email_notification",
|
||||
python_callable=notify_email,
|
||||
provide_context=True,
|
||||
dag=dag,
|
||||
trigger_rule=TriggerRule.ALL_DONE,
|
||||
op_args=[LOCAL_CALCULATION_DT, DAG_RUN_BEGIN_TIME],
|
||||
)
|
||||
|
||||
... >> email_notification_task
|
||||
```
|
||||
|
||||
## Практическое применение: улучшенный DAG
|
||||
|
||||
Теперь применим изученные концепции для усовершенствования нашего DAG. Добавим переменные и разобьем задачи на логические группы.
|
||||
|
||||
Сначала создадим переменную `DATABASE_URL` со строкой подключения к базе данных через веб-интерфейс Airflow и импортируем ее в коде DAG:
|
||||
|
||||
В этом примере используется переменная 'database_connection_string', предварительно созданная в интерфейсе Airflow, для хранения строки подключения к базе данных. Это позволяет избежать жесткого кодирования конфиденциальной информации в коде DAG и упрощает настройку подключения для разных окружений.
|
||||
|
||||
```python
|
||||
import datetime as dt
|
||||
import pandas as pd
|
||||
from airflow.models import DAG
|
||||
from airflow.operators.bash import BashOperator
|
||||
from airflow.operators.python import PythonOperator
|
||||
from airflow.operators.dummy import DummyOperator
|
||||
from airflow.utils.task_group import TaskGroup
|
||||
from airflow.models import Variable
|
||||
from sqlalchemy import create_engine
|
||||
|
||||
DATABASE_URL = Variable.get('database_connection_string')
|
||||
|
||||
args = {
|
||||
'owner': 'analytics_team',
|
||||
'start_date': dt.datetime(2021, 6, 15),
|
||||
'retries': 2,
|
||||
'retry_delay': dt.timedelta(minutes=2),
|
||||
}
|
||||
|
||||
# функции для обработки данных
|
||||
def get_file_path(file_name):
|
||||
return os.path.join(os.path.expanduser('~/data'), file_name)
|
||||
|
||||
def load_customer_data():
|
||||
url = 'https://example.com/customer_data.csv'
|
||||
df = pd.read_csv(url)
|
||||
engine = create_engine(DATABASE_URL)
|
||||
df.to_sql('customers', engine, index=False, if_exists='replace', schema='staging')
|
||||
|
||||
def aggregate_customer_data():
|
||||
engine = create_engine(DATABASE_URL)
|
||||
customer_df = pd.read_sql('select * from staging.customers', con=engine)
|
||||
|
||||
df = customer_df.groupby(['region', 'category']).agg(
|
||||
total_orders=('orders', 'sum'),
|
||||
avg_amount=('amount', 'mean')
|
||||
).reset_index()
|
||||
|
||||
df.to_sql('customer_summary', engine, index=False, if_exists='replace', schema='analytics')
|
||||
|
||||
dag = DAG(
|
||||
dag_id='customer_pipeline_enhanced',
|
||||
schedule_interval=None,
|
||||
default_args=args,
|
||||
)
|
||||
```
|
||||
|
||||
Теперь разобьем задачи на логические группы и добавим Jinja-шаблоны для доступа к контексту выполнения:
|
||||
|
||||
В этом примере создается начальная задача 'pipeline_start', которая использует Jinja-шаблоны для вывода информации о запуске DAG, включая идентификатор запуска (run_id) и информацию о DAG Run. Затем задачи группируются в логическую группу 'data_processing_stage', что улучшает структуру и читаемость DAG.
|
||||
|
||||
```python
|
||||
# Начальная задача с информацией о запуске
|
||||
start_task = BashOperator(
|
||||
task_id='pipeline_start',
|
||||
bash_command='echo "Pipeline started! Run ID: {{ run_id }} | DAG Run: {{ dag_run }}"',
|
||||
dag=dag,
|
||||
)
|
||||
|
||||
# Группа задач по предварительной обработке данных
|
||||
with TaskGroup(group_id="data_processing_stage") as data_processing:
|
||||
# Загрузка данных
|
||||
load_customer_dataset = PythonOperator(
|
||||
task_id='load_customer_data',
|
||||
python_callable=load_customer_data,
|
||||
dag=dag,
|
||||
)
|
||||
# Агрегация и запись данных
|
||||
aggregate_customer_dataset = PythonOperator(
|
||||
task_id='aggregate_customer_data',
|
||||
python_callable=aggregate_customer_data,
|
||||
dag=dag,
|
||||
)
|
||||
load_customer_dataset >> aggregate_customer_dataset
|
||||
|
||||
# Установка последовательности выполнения
|
||||
start_task >> data_processing
|
||||
```
|
||||
|
||||
Поскольку последовательность задач внутри групп указывается при их создании, в конце необходимо определить порядок выполнения самих групп, чтобы планировщик понимал общую логику выполнения.
|
||||
|
||||
Итоговый DAG будет выглядеть следующим образом:
|
||||
|
||||
```python
|
||||
import datetime as dt
|
||||
import pandas as pd
|
||||
from airflow.models import DAG
|
||||
from airflow.operators.bash import BashOperator
|
||||
from airflow.operators.python import PythonOperator
|
||||
from airflow.operators.dummy import DummyOperator
|
||||
from airflow.utils.task_group import TaskGroup
|
||||
from airflow.models import Variable
|
||||
from sqlalchemy import create_engine
|
||||
|
||||
DATABASE_URL = Variable.get('database_connection_string')
|
||||
|
||||
args = {
|
||||
'owner': 'analytics_team',
|
||||
'start_date': dt.datetime(2021, 6, 15),
|
||||
'retries': 2,
|
||||
'retry_delay': dt.timedelta(minutes=2),
|
||||
}
|
||||
|
||||
def get_file_path(file_name):
|
||||
return os.path.join(os.path.expanduser('~/data'), file_name)
|
||||
|
||||
def load_customer_data():
|
||||
url = 'https://example.com/customer_data.csv'
|
||||
df = pd.read_csv(url)
|
||||
engine = create_engine(DATABASE_URL)
|
||||
df.to_sql('customers', engine, index=False, if_exists='replace', schema='staging')
|
||||
|
||||
def aggregate_customer_data():
|
||||
engine = create_engine(DATABASE_URL)
|
||||
customer_df = pd.read_sql('select * from staging.customers', con=engine)
|
||||
|
||||
df = customer_df.groupby(['region', 'category']).agg(
|
||||
total_orders=('orders', 'sum'),
|
||||
avg_amount=('amount', 'mean')
|
||||
).reset_index()
|
||||
|
||||
df.to_sql('customer_summary', engine, index=False, if_exists='replace', schema='analytics')
|
||||
|
||||
dag = DAG(
|
||||
dag_id='customer_pipeline_enhanced',
|
||||
schedule_interval=None,
|
||||
default_args=args,
|
||||
)
|
||||
|
||||
# Начальная задача
|
||||
start_task = BashOperator(
|
||||
task_id='pipeline_start',
|
||||
bash_command='echo "Pipeline started! Run ID: {{ run_id }} | DAG Run: {{ dag_run }}"',
|
||||
dag=dag,
|
||||
)
|
||||
|
||||
# Группа предварительной обработки
|
||||
with TaskGroup(group_id="data_processing_stage") as data_processing:
|
||||
load_customer_dataset = PythonOperator(
|
||||
task_id='load_customer_data',
|
||||
python_callable=load_customer_data,
|
||||
dag=dag,
|
||||
)
|
||||
aggregate_customer_dataset = PythonOperator(
|
||||
task_id='aggregate_customer_data',
|
||||
python_callable=aggregate_customer_data,
|
||||
dag=dag,
|
||||
)
|
||||
load_customer_dataset >> aggregate_customer_dataset
|
||||
|
||||
start_task >> data_processing
|
||||
```
|
||||
|
||||
В этом материале мы рассмотрели расширенные возможности Airflow, которые помогут улучшить работу ваших пайплайнов:
|
||||
- Управление ресурсами с помощью пулов задач
|
||||
- Обмен данными между задачами через XCom
|
||||
- Логическая группировка задач с TaskGroup
|
||||
- Настройка системы оповещений
|
||||
|
||||
Помните: не стоит использовать все доступные функции сразу. Выбирайте инструменты последовательно и находите оптимальный набор возможностей под конкретную задачу.
|
||||
@@ -0,0 +1,98 @@
|
||||
# Учебник по Apache Airflow для начинающих
|
||||
|
||||
Добро пожаловать в учебник по Apache Airflow! Этот курс создан специально для тех, кто уже освоил Python и SQL и готов познакомиться со своим первым инструментом для ETL-процессов.
|
||||
|
||||
## 📚 Введение
|
||||
|
||||
Apache Airflow — это мощный оркестратор рабочих процессов с открытым исходным кодом, который стал стандартом де-факто в мире обработки данных. Если вы работаете с автоматизацией процессов обработки информации — будь то подготовка данных для машинного обучения, создание отчетов или интеграция данных из различных источников — Airflow поможет вам организовать эти процессы эффективно и надежно.
|
||||
|
||||
Этот учебник разработан с учетом того, что вы:
|
||||
- Имеете базовые знания Python и SQL
|
||||
- Хотите освоить свой первый ETL-инструмент
|
||||
- Нуждаетесь в практическом руководстве с понятными объяснениями
|
||||
|
||||
Мы начнем с основ и постепенно перейдем к продвинутым возможностям, чтобы вы могли уверенно использовать Airflow в реальных проектах. Курс ориентирован на версию Airflow 2.5 и использует практический подход с множеством примеров и визуальных материалов.
|
||||
|
||||
## 🎯 Что вы узнаете
|
||||
|
||||
- **Основы Airflow**: что это такое, для каких задач используется и почему он стал таким популярным
|
||||
- **Ключевые концепции**: DAG, задачи, операторы и другие фундаментальные понятия
|
||||
- **Архитектура**: как устроен Airflow внутри и как это влияет на вашу работу
|
||||
- **Интерфейс**: как эффективно использовать веб-интерфейс для мониторинга и управления
|
||||
- **Практическое создание DAG**: пошаговое руководство по созданию ваших первых рабочих процессов
|
||||
- **Отладка и мониторинг**: как понимать статусы задач и быстро находить проблемы
|
||||
- **Продвинутые возможности**: шаблоны, параметры, управление временем и другие мощные функции
|
||||
|
||||
## 📖 Содержание курса
|
||||
|
||||
### [01. Введение в Airflow](01%20-%20Введение%20в%20Airflow.md)
|
||||
- Почему Apache Airflow стал незаменимым инструментом для работы с данными
|
||||
- Преимущества использования Airflow
|
||||
- Когда Airflow может не подойти
|
||||
- Практическое применение Airflow
|
||||
- Как компании развертывают Airflow
|
||||
|
||||
### [02. Понятное введение в ключевые понятия Airflow](02%20-%20Понятное%20введение%20в%20ключевые%20понятия%20Airflow.md)
|
||||
- Основные концепции и терминология
|
||||
- Что такое DAG и как он работает
|
||||
- Операторы, задачи и зависимости
|
||||
|
||||
### [03. Как устроен Airflow внутри](03%20-%20Как%20устроен%20Airflow%20внутри.md)
|
||||
- Архитектура Airflow
|
||||
- Компоненты системы
|
||||
- Как работает планировщик и исполнитель
|
||||
|
||||
### [04. Знакомство с интерфейсом Airflow для начинающих](04%20-%20Знакомство%20с%20интерфейсом%20Airflow%20для%20начинающих.md)
|
||||
- Навигация по веб-интерфейсу
|
||||
- Основные разделы и их назначение
|
||||
- Как мониторить и управлять DAG
|
||||
|
||||
### [05. Основы построения DAG-файлов в Airflow](05%20-%20Основы%20построения%20DAG-файлов%20в%20Airflow.md)
|
||||
- Структура DAG-файла
|
||||
- Создание ваших первых задач
|
||||
- Настройка расписаний и параметров
|
||||
|
||||
### [06. Статусы задач в Airflow](06%20-%20Статусы%20задач%20в%20Airflow.md)
|
||||
- Что означают цвета и статусы задач
|
||||
- Жизненный цикл задачи
|
||||
- Как отлаживать проблемы и работать с ошибками
|
||||
|
||||
### [07. Гибкие шаблоны и настройки в Airflow](07%20-%20Гибкие%20шаблоны%20и%20настройки%20в%20Airflow.md)
|
||||
- Использование Jinja-шаблонов
|
||||
- Параметризация DAG
|
||||
- Глобальные переменные и соединения
|
||||
|
||||
### [08. Управление временем в Airflow](08%20-%20Управление%20временем.md)
|
||||
- Расписания и интервалы запуска
|
||||
- Работа с временными зонами
|
||||
- Понимание execution_date и других временных концепций
|
||||
|
||||
### [09. Продвинутые возможности Airflow](09%20-%20Продвинутые%20возможности%20Airflow.md)
|
||||
- Расширенные операторы и сенсоры
|
||||
- Обработка ошибок и повторные попытки
|
||||
- Оптимизация и масштабирование DAG
|
||||
|
||||
## 🚀 Как использовать этот учебник
|
||||
|
||||
1. **Следуйте порядку**: начните с Введения в Airflow и переходите к следующим по порядку
|
||||
2. **Практикуйтесь**: создавайте свои DAG по мере изучения материала
|
||||
3. **Используйте изображения**: в учебнике много визуальных материалов для лучшего понимания
|
||||
4. **Экспериментируйте**: не бойтесь пробовать разные настройки и сценарии
|
||||
|
||||
## 📝 Технические требования
|
||||
|
||||
- **Python**: базовое знание языка
|
||||
- **SQL**: понимание основных запросов и концепций
|
||||
- **Airflow**: версия 2.5 (на которой ориентирован курс)
|
||||
- **Среда разработки**: любой текстовый редактор или IDE с поддержкой Python
|
||||
|
||||
## 💡 Советы для успешного обучения
|
||||
|
||||
- **Не торопитесь**: лучше глубоко понять каждый концепт, чем быстро пройти весь курс
|
||||
- **Задавайте вопросы**: если что-то непонятно, вернитесь к предыдущим материалам
|
||||
- **Делайте заметки**: записывайте важные моменты и примеры кода
|
||||
- **Применяйте на практике**: сразу пробуйте создавать свои DAG для реальных задач
|
||||
|
||||
---
|
||||
|
||||
*Этот учебник создан для начинающих специалистов в области данных и инженерии. Все материалы ориентированы на практическое применение и пошаговое освоение Apache Airflow.*
|
||||
|
After Width: | Height: | Size: 52 KiB |
|
After Width: | Height: | Size: 1.3 MiB |
|
After Width: | Height: | Size: 1.4 MiB |
|
After Width: | Height: | Size: 9.6 MiB |
|
After Width: | Height: | Size: 89 KiB |
|
After Width: | Height: | Size: 160 KiB |
|
After Width: | Height: | Size: 25 KiB |
|
After Width: | Height: | Size: 60 KiB |
|
After Width: | Height: | Size: 28 KiB |
|
After Width: | Height: | Size: 778 KiB |
|
After Width: | Height: | Size: 1.0 MiB |
|
After Width: | Height: | Size: 34 KiB |
|
After Width: | Height: | Size: 475 KiB |
|
After Width: | Height: | Size: 187 KiB |
|
After Width: | Height: | Size: 311 KiB |
|
After Width: | Height: | Size: 807 KiB |
|
After Width: | Height: | Size: 838 KiB |
|
After Width: | Height: | Size: 880 KiB |
|
After Width: | Height: | Size: 873 KiB |
|
After Width: | Height: | Size: 322 KiB |
|
After Width: | Height: | Size: 834 KiB |
|
After Width: | Height: | Size: 142 KiB |
|
After Width: | Height: | Size: 188 KiB |
|
After Width: | Height: | Size: 649 KiB |
|
After Width: | Height: | Size: 2.0 MiB |
|
After Width: | Height: | Size: 502 KiB |
|
After Width: | Height: | Size: 38 KiB |
|
After Width: | Height: | Size: 164 KiB |
|
After Width: | Height: | Size: 22 KiB |
|
After Width: | Height: | Size: 546 KiB |
|
After Width: | Height: | Size: 748 KiB |
|
After Width: | Height: | Size: 2.4 MiB |