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

|
```mermaid
|
||||||
|
flowchart LR
|
||||||
|
in[Данные] --> op1[Операция № 1] --> op2[Операция № 2] --> op3[Операция № 3] --> out[Данные]
|
||||||
|
|
||||||
|
subgraph PIPE[Конвейер]
|
||||||
|
op1
|
||||||
|
op2
|
||||||
|
op3
|
||||||
|
end
|
||||||
|
```
|
||||||
|
|
||||||
Аналогичная ситуация возникает при обработке данных: в определенный момент необходимо собрать актуальную информацию, преобразовать её и выполнить различные операции — создать витрину данных и сохранить в базу, обучить модель машинного обучения, подготовить отчет в Excel и разослать его по электронной почте. Вариантов множество, и для решения таких задач требуются специализированные инструменты.
|
Аналогичная ситуация возникает при обработке данных: в определенный момент необходимо собрать актуальную информацию, преобразовать её и выполнить различные операции — создать витрину данных и сохранить в базу, обучить модель машинного обучения, подготовить отчет в Excel и разослать его по электронной почте. Вариантов множество, и для решения таких задач требуются специализированные инструменты.
|
||||||
|
|
||||||
@@ -77,4 +85,4 @@
|
|||||||
|
|
||||||
В этом материале вы познакомились с предпосылками появления инструментов управления процессами, подобных Airflow. Вы узнали, в каких сценариях Airflow наиболее эффективен, а в каких случаях стоит рассмотреть альтернативные решения, а также как этот инструмент применяется на практике в крупных российских и международных компаниях.
|
В этом материале вы познакомились с предпосылками появления инструментов управления процессами, подобных Airflow. Вы узнали, в каких сценариях Airflow наиболее эффективен, а в каких случаях стоит рассмотреть альтернативные решения, а также как этот инструмент применяется на практике в крупных российских и международных компаниях.
|
||||||
|
|
||||||
💡 Обратите внимание, что в данном модуле рассматривается Airflow версии 2.5.
|
💡 Обратите внимание, что в данном модуле рассматривается Airflow версии 2.9.
|
||||||
|
|||||||
@@ -12,7 +12,18 @@
|
|||||||
|
|
||||||
Проще говоря, DAG — это упорядоченный набор задач, которые выполняются строго по расписанию и никогда не повторяются в рамках одного запуска. Создавая пайплайны для обработки данных, вы фактически создаете DAG.
|
Проще говоря, DAG — это упорядоченный набор задач, которые выполняются строго по расписанию и никогда не повторяются в рамках одного запуска. Создавая пайплайны для обработки данных, вы фактически создаете DAG.
|
||||||
|
|
||||||

|
```mermaid
|
||||||
|
flowchart LR
|
||||||
|
n1((1)) --> n2((2))
|
||||||
|
n1 --> n3((3))
|
||||||
|
|
||||||
|
n2 --> n4((4))
|
||||||
|
n3 --> n6((6))
|
||||||
|
n6 --> n7((7))
|
||||||
|
|
||||||
|
n7 --> n12
|
||||||
|
n4 --> n12((12))
|
||||||
|
```
|
||||||
|
|
||||||
# Шаги вашего пайплайна: Задачи
|
# Шаги вашего пайплайна: Задачи
|
||||||
|
|
||||||
@@ -26,14 +37,52 @@
|
|||||||
|
|
||||||
На примере ниже видно, как задача E зависит от успешного завершения всех предыдущих задач:
|
На примере ниже видно, как задача E зависит от успешного завершения всех предыдущих задач:
|
||||||
|
|
||||||

|
```mermaid
|
||||||
|
flowchart LR
|
||||||
|
A["Task A"] --> B["Task B"]
|
||||||
|
A --> C["Task C"]
|
||||||
|
|
||||||
|
B --> D["Task D"]
|
||||||
|
C --> D
|
||||||
|
|
||||||
|
D --> E["Task E"]
|
||||||
|
```
|
||||||
|
|
||||||
Airflow позволяет создавать сложные сценарии:
|
Airflow позволяет создавать сложные сценарии:
|
||||||
- Зависимости между разными DAG-ами с помощью TriggerDagRunOperator и ExternalTaskSensor
|
- Зависимости между разными DAG-ами с помощью TriggerDagRunOperator и ExternalTaskSensor
|
||||||
- Условное выполнение задач в зависимости от результатов предыдущих шагов
|
- Условное выполнение задач в зависимости от результатов предыдущих шагов
|
||||||
- Сложные ветвления и параллельные ветки выполнения
|
- Сложные ветвления и параллельные ветки выполнения
|
||||||
|
|
||||||

|
```mermaid
|
||||||
|
flowchart LR
|
||||||
|
%% без цветов, только форма и рамка
|
||||||
|
classDef taskGroup stroke-width:3px;
|
||||||
|
classDef endpoint stroke-width:2px,stroke-dasharray: 5 3;
|
||||||
|
|
||||||
|
start(("начало")):::endpoint --> print_start["print_start_bash"]
|
||||||
|
|
||||||
|
print_start --> py1["python_function_with_input_1"]
|
||||||
|
print_start --> py4["python_function_with_input_4"]
|
||||||
|
|
||||||
|
py1 --> py2["python_function_with_input_2"]
|
||||||
|
py1 --> py3["python_function_with_input_3"]
|
||||||
|
|
||||||
|
py4 --> py5["python_function_with_input_5"]
|
||||||
|
py4 --> py6["python_function_with_input_6"]
|
||||||
|
|
||||||
|
py2 --> py_join["python_function_with_input"]
|
||||||
|
py3 --> py_join
|
||||||
|
py5 --> py_join
|
||||||
|
py6 --> py_join
|
||||||
|
|
||||||
|
py_join --> tg["task_group_with_two_tasks"]:::taskGroup
|
||||||
|
|
||||||
|
tg --> dyn123["dynamic_task_123"]
|
||||||
|
tg --> dyn456["dynamic_task_456"]
|
||||||
|
|
||||||
|
dyn123 --> finish(("конец")):::endpoint
|
||||||
|
dyn456 --> finish
|
||||||
|
```
|
||||||
|
|
||||||
# Инструменты для выполнения: Операторы
|
# Инструменты для выполнения: Операторы
|
||||||
|
|
||||||
|
|||||||
@@ -8,7 +8,36 @@ Apache Airflow состоит из нескольких взаимосвязан
|
|||||||
|
|
||||||
На схеме ниже показано, как эти компоненты взаимодействуют между собой:
|
На схеме ниже показано, как эти компоненты взаимодействуют между собой:
|
||||||
|
|
||||||

|
```mermaid
|
||||||
|
flowchart LR
|
||||||
|
de["👤 Data Engineer"]
|
||||||
|
dags["DAGs (файлы)"]
|
||||||
|
cfg["airflow.cfg"]
|
||||||
|
|
||||||
|
subgraph core["Airflow core"]
|
||||||
|
web["Webserver (UI)"]
|
||||||
|
|
||||||
|
subgraph se["Scheduling & Execution"]
|
||||||
|
sch["Scheduler"]
|
||||||
|
ex["Executor"]
|
||||||
|
wk["Worker(s)"]
|
||||||
|
|
||||||
|
sch --> ex
|
||||||
|
ex --> wk
|
||||||
|
end
|
||||||
|
end
|
||||||
|
|
||||||
|
db["Metadata DB (Postgres)"]
|
||||||
|
|
||||||
|
de --> web
|
||||||
|
de --> dags
|
||||||
|
de --> cfg
|
||||||
|
|
||||||
|
dags --> sch
|
||||||
|
|
||||||
|
web --> db
|
||||||
|
se --> db
|
||||||
|
```
|
||||||
*Как устроена система Airflow*
|
*Как устроена система Airflow*
|
||||||
|
|
||||||
Давайте подробно рассмотрим каждый компонент системы.
|
Давайте подробно рассмотрим каждый компонент системы.
|
||||||
|
|||||||
@@ -1,114 +1,88 @@
|
|||||||
# Знакомство с интерфейсом Airflow для начинающих
|
# Знакомство с интерфейсом Airflow для начинающих
|
||||||
|
|
||||||
Когда вы создаете ETL-процессы для автоматической обработки данных, важно иметь удобный способ отслеживать их работу, находить ошибки и управлять выполнением. Именно для этого в Apache Airflow предусмотрен веб-интерфейс — ваш главный помощник в повседневной работе с пайплайнами.
|
В предыдущих главах вы познакомились с основными понятиями Apache Airflow и поняли, из каких компонентов он состоит внутри. Теперь самое время открыть веб‑интерфейс — главный инструмент для ежедневной работы с пайплайнами.
|
||||||
|
|
||||||
В этом материале мы подробно разберем, как устроен интерфейс Airflow версии 2.5 и какие возможности он предоставляет для мониторинга и управления вашими процессами обработки данных.
|
Через UI вы будете:
|
||||||
|
|
||||||
# Домашняя страница Airflow
|
- смотреть список DAG-ов и их статусы;
|
||||||
|
- запускать пайплайны вручную и останавливать их;
|
||||||
|
- разбираться, почему задачи упали, и перезапускать их;
|
||||||
|
- анализировать время выполнения и находить узкие места.
|
||||||
|
|
||||||
После входа в систему вы попадете на главную страницу — центральную панель управления, где собрана вся ключевая информация о ваших пайплайнах. Не пугайтесь обилия элементов и цветов — все устроено логично и интуитивно понятно.
|
В этой главе мы сосредоточимся на ключевых экранах интерфейса Airflow 2.9.x и будем опираться на стенд из папки `airflow-docker` в этом репозитории.
|
||||||
|
|
||||||
По сути, это обычная таблица, где каждая строка представляет собой один DAG (Directed Acyclic Graph) — ваш пайплайн обработки данных, а столбцы содержат различную информацию о нем.
|
## Главная страница: список DAG-ов (DAGs View)
|
||||||
|
|
||||||
## Основные элементы главной страницы
|
После авторизации в Airflow вы попадаете на главную страницу — список всех DAG-ов. Это экран, откуда начинается почти любая работа.
|
||||||
|
|
||||||
**Название DAG** — в первом столбце отображается список всех зарегистрированных в системе пайплайнов. По умолчанию Airflow включает демонстрационные примеры различных операторов. Список отсортирован по алфавиту для удобства поиска.
|
Здесь полезно знать несколько ключевых колонок (названия — как в англоязычном UI):
|
||||||
|
|
||||||

|
- **DAG / DAG ID** — идентификатор пайплайна; клик по нему открывает страницу конкретного DAG-а.
|
||||||
|
- **Owner** — владелец пайплайна (на кого «вешать» вопросы и инциденты).
|
||||||
|
- **Schedule** — расписание запуска в виде cron-выражения или пресета.
|
||||||
|
- **Last Run / Next Run** — когда DAG запускался в последний раз и когда запустится по расписанию.
|
||||||
|
- **Recent Tasks** — короткая сводка по результатам последних запусков (сколько задач `success` / `failed` / `running` и т.п.).
|
||||||
|
- **Actions** — быстрые действия: переключатель **Pause / Unpause** (включить/выключить выполнение по расписанию), кнопка **Trigger DAG** (ручной запуск), ссылки на детальные представления: **Grid**, **Graph**, иногда **Code**.
|
||||||
|
|
||||||
**Переключатели активности** — напротив каждого DAG находится кнопка-выключатель, позволяющая мгновенно активировать или деактивировать пайплайн прямо из веб-интерфейса без изменения кода.
|
На этом экране вы решаете простой вопрос:
|
||||||
|
«С моими DAG-ами всё более-менее нормально или где-то горит?»
|
||||||
|
|
||||||

|
## Страница DAG: главное рабочее место
|
||||||
|
|
||||||
**Владелец процесса** — каждый пайплайн имеет ответственного владельца. Это особенно полезно в командной работе, когда несколько инженеров создают и поддерживают различные ETL-процессы. Владелец отвечает за мониторинг и корректную работу своего DAG.
|
Когда на главной странице (**DAGs View**) вы нажимаете на `DAG ID`, открывается страница конкретного пайплайна.
|
||||||
|
|
||||||

|
Верхняя часть страницы DAG:
|
||||||
|
|
||||||
**Статус выполнения** — цветные индикаторы с цифрами показывают количество и состояние последних запусков DAG:
|
* переключатель **Pause / Unpause**;
|
||||||
- 🔴 Красный — завершено с ошибкой (failed)
|
* кнопка **Trigger DAG** (ручной запуск);
|
||||||
- 🟡 Желтый — ожидает повторного запуска (retry)
|
* фильтр по дате, типу и состоянию запусков (**Run Type**, **Run State**, период по календарю);
|
||||||
- 🟢 Зеленый — выполняется в данный момент (running)
|
* небольшой индикатор статусов задач (цветные ярлыки `running`, `failed`, `success` и т.п.).
|
||||||
- 🟢 Тёмно-зеленый — успешно завершено (success)
|
|
||||||
|
|
||||||

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

|
* **Details**
|
||||||
|
Краткое резюме DAG: количество задач, типы операторов, расписание, теги, статистика по запускам. Это удобная точка входа: «что это за DAG и как он в целом живёт».
|
||||||
|
|
||||||
**Последний запуск** — показывает дату и время самого свежего выполнения DAG, будь то автоматический запуск по расписанию или ручной запуск.
|
* **Graph**
|
||||||
|
Граф зависимостей задач. Здесь хорошо видно, какие задачи идут последовательно, какие — параллельно, где ветвления.
|
||||||
|
Клик по задаче открывает панель с действиями: **View Log**, **Clear**, **Mark Success / Mark Failed**, **Run** и др.
|
||||||
|
|
||||||

|
* **Gantt**
|
||||||
|
Диаграмма Ганта для выбранного запуска DAG. Показывает, сколько времени заняла каждая задача и где они выполнялись параллельно. По ней удобно искать «бутылочные горлышки» — самые долгие шаги пайплайна.
|
||||||
|
|
||||||
**Статус задач** — детальная информация о последнем запуске: сколько задач находится в каждом статусе. Это помогает быстро оценить общее состояние пайплайна без необходимости погружаться в детали.
|
* **Run Duration**
|
||||||
|
История длительности запусков DAG. Помогает увидеть, не стали ли запуски в целом работать заметно дольше, и отследить, после какого изменения время выполнения выросло.
|
||||||
|
|
||||||

|
* **Calendar**
|
||||||
|
Календарный вид истории запусков: по дням и месяцам видно, когда DAG запускался и как часто были ошибки.
|
||||||
|
|
||||||
**Быстрые действия** — в последнем столбце расположены кнопки для немедленного выполнения операций: запуск, обновление и удаление DAG. На практике этими кнопками пользуются редко.
|
* **Code**
|
||||||
|
Исходный код DAG, который сейчас задеплоен в Airflow. Быстрый способ проверить, что в среде действительно лежит та версия DAG, которую вы ждёте (и что изменения из Git уже подхватились).
|
||||||
|
|
||||||
Главная страница дает вам общее представление о состоянии всех ваших процессов. Но для детальной работы с конкретным пайплайном нужно перейти внутрь — просто кликните по названию интересующего DAG.
|
* **Audit Log**
|
||||||
|
Журнал действий по DAG: кто запускал, очищал задачи, менял состояние и т.д. Полезен, когда нужно понять, «кто и что нажал» перед тем, как всё сломалось.
|
||||||
|
|
||||||

|
### Как работать с задачами (Tasks)
|
||||||
|
|
||||||
# Страница конкретного DAG
|
Независимо от вкладки (чаще всего — **Graph** или **Gantt**), логика одна:
|
||||||
|
|
||||||
После перехода внутрь DAG вы увидите набор вкладок с различной информацией: от визуального представления структуры пайплайна до детальных логов выполнения и исходного кода.
|
1. Находите нужную задачу.
|
||||||
|
2. Кликаете по ней — справа (или во всплывающем окне) появляется панель **Task Instance**.
|
||||||
|
3. В этой панели доступны:
|
||||||
|
|
||||||
## Древовидное представление (Tree View)
|
* **View Log** — открыть логи;
|
||||||
|
* **Clear** — очистить состояние для повторного запуска;
|
||||||
|
* **Mark Success / Mark Failed** — вручную выставить статус;
|
||||||
|
* **Run** — запустить задачу сейчас.
|
||||||
|
|
||||||
По умолчанию открывается вкладка с древовидной структурой задач. Здесь отображаются все запуски DAG с указанием статуса каждой задачи, времени выполнения и других метрик мониторинга.
|
Через дополнительные опции **Clear** можно захватывать **upstream** / **downstream** задачи и несколько запусков сразу — это основной инструмент «перезапуска кусочка» DAG.
|
||||||
|
|
||||||
Вы можете увидеть:
|
Подробное описание всех экранов Airflow (с актуальными скриншотами) есть в официальной документации: [UI / Screenshots (Apache Airflow 2.9.3)][1].
|
||||||
- Состав DAG и последовательность выполнения задач
|
|
||||||
- Тип оператора для каждой задачи
|
|
||||||
- Историю запусков в виде цветных квадратов напротив каждой задачи
|
|
||||||
|
|
||||||

|
Там же описаны **DAGs View**, **Grid View**, **Graph View**, **Gantt Chart**, **Task Duration**, **Landing Times**, **Code View**, **Audit Log** и другие разделы UI.
|
||||||

|
|
||||||
|
|
||||||
## Графическое представление (Graph View)
|
---
|
||||||
|
|
||||||
Когда DAG содержит много задач, древовидное представление может быть неудобным. В таких случаях используйте вкладку Graph View — она показывает пайплайн в виде наглядного графа с четкими связями между задачами.
|
[1]: https://airflow.apache.org/docs/apache-airflow/2.9.3/ui.html "UI / Screenshots — Airflow Documentation"
|
||||||
|
|
||||||

|
|
||||||

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

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

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

|
|
||||||
|
|
||||||
Теперь вы знакомы с основными возможностями веб-интерфейса Airflow для мониторинга и управления вашими процессами обработки данных. Эти знания помогут вам эффективно работать с пайплайнами и быстро решать возникающие проблемы.
|
|
||||||
|
|||||||
@@ -4,20 +4,63 @@
|
|||||||
|
|
||||||
# Как устроен код DAG-файла
|
# Как устроен код DAG-файла
|
||||||
|
|
||||||
Как вы уже знаете из предыдущих уроков, DAG представляет собой граф вычислений, состоящий из отдельных задач (tasks). На уровне кода DAG — это обычный Python-файл, который описывает все задачи в рамках пайплайна и определяет последовательность их выполнения.
|
DAG представляет собой граф вычислений, состоящий из отдельных задач (tasks). На уровне кода DAG‑файл — это обычный Python‑скрипт, который Airflow регулярно импортирует, чтобы «увидеть» ваши пайплайны.
|
||||||
|
|
||||||
Любой DAG-файл состоит из нескольких ключевых компонентов:
|
Удобно мысленно разбивать любой DAG‑файл на четыре блока:
|
||||||
- Импорт необходимых модулей и библиотек
|
1. Импорт необходимых модулей и операторов
|
||||||
- Настройка параметров и инициализация объекта DAG
|
2. Настройка параметров и инициализация объекта `DAG`
|
||||||
- Создание отдельных задач с помощью операторов
|
3. Создание отдельных задач с помощью операторов
|
||||||
- Определение порядка выполнения задач
|
4. Определение порядка выполнения задач (задание зависимостей)
|
||||||
|
|
||||||
Порядок этих компонентов имеет значение, поскольку Airflow — это Python-библиотека, и обращение к еще не инициализированным объектам приведет к ошибкам выполнения.
|
Порядок этих блоков важен: как и в любом Python‑коде, нельзя обращаться к объектам, которые ещё не созданы.
|
||||||
|
|
||||||
Давайте рассмотрим простой пример DAG-файла и разберем его по частям.
|
Ниже приведён простой пример DAG‑файла. В следующих подразделах мы разберём каждый блок по отдельности, опираясь на этот пример.
|
||||||
|
|
||||||

|
```python
|
||||||
*Пример структуры DAG-файла*
|
# Секция импортов
|
||||||
|
from datetime import timedelta
|
||||||
|
from airflow import DAG
|
||||||
|
from airflow.operators.bash import BashOperator
|
||||||
|
from airflow.utils.dates import days_ago
|
||||||
|
|
||||||
|
# Настройки параметров по умолчанию (default_args)
|
||||||
|
default_args = {
|
||||||
|
'owner': 'airflow',
|
||||||
|
'depends_on_past': False,
|
||||||
|
'start_date': days_ago(2),
|
||||||
|
'email': ['airflow@example.com'],
|
||||||
|
'email_on_failure': False,
|
||||||
|
'email_on_retry': False,
|
||||||
|
'retries': 1,
|
||||||
|
'retry_delay': timedelta(minutes=5),
|
||||||
|
}
|
||||||
|
|
||||||
|
# Объявление DAG
|
||||||
|
dag = DAG(
|
||||||
|
dag_id='tutorial',
|
||||||
|
default_args=default_args,
|
||||||
|
description='A simple tutorial DAG',
|
||||||
|
schedule_interval=timedelta(days=1),
|
||||||
|
)
|
||||||
|
|
||||||
|
# Объявление задач (operators)
|
||||||
|
t1 = BashOperator(
|
||||||
|
task_id='print_date',
|
||||||
|
bash_command='date',
|
||||||
|
dag=dag,
|
||||||
|
)
|
||||||
|
|
||||||
|
t2 = BashOperator(
|
||||||
|
task_id='sleep',
|
||||||
|
depends_on_past=False,
|
||||||
|
bash_command='sleep 5',
|
||||||
|
retries=3,
|
||||||
|
dag=dag,
|
||||||
|
)
|
||||||
|
|
||||||
|
# Задание графа последовательности
|
||||||
|
t1 >> t2
|
||||||
|
```
|
||||||
|
|
||||||
## Импорт необходимых модулей
|
## Импорт необходимых модулей
|
||||||
|
|
||||||
@@ -90,35 +133,57 @@ t2 = BashOperator(
|
|||||||
|
|
||||||
## Определение последовательности выполнения
|
## Определение последовательности выполнения
|
||||||
|
|
||||||
Завершающий этап — указание порядка выполнения задач. В Airflow для этого используются стрелочные операторы:
|
Завершающий этап — указание порядка выполнения задач. В Airflow для этого используются стрелочные операторы `>>` и `<<`.
|
||||||
|
|
||||||
В приведенном примере задача с идентификатором 'process_data' будет запускаться только после успешного завершения задачи 'show_time'.
|
В простейшем случае мы просто строим цепочку:
|
||||||
|
|
||||||
```python
|
```python
|
||||||
t1 >> t2 # t2 запускается после завершения t1
|
t1 >> t2 >> t3 # t2 после t1, t3 после t2
|
||||||
```
|
```
|
||||||
|
|
||||||
Альтернативный синтаксис:
|
Параллельные ветки:
|
||||||
|
|
||||||
Этот синтаксис эквивалентен предыдущему примеру, просто записан в обратном порядке. Задача 'show_time' должна завершиться перед запуском задачи 'process_data'.
|
|
||||||
```python
|
|
||||||
t2 << t1 # t1 должна завершиться перед запуском t2
|
|
||||||
```
|
|
||||||
|
|
||||||
Можно также группировать задачи в списки для создания более сложных зависимостей:
|
|
||||||
|
|
||||||
В этом примере задачи 't2' и 't3' будут запускаться одновременно после завершения задачи 't1'. Затем задача 't4' запустится после завершения обеих задач 't2' и 't3'.
|
|
||||||
|
|
||||||
```python
|
```python
|
||||||
t1 >> [t2, t3] # t2 и t3 запускаются одновременно после t1
|
t1 >> [t2, t3] # t2 и t3 стартуют после t1
|
||||||
[t2, t3] >> t4 # t4 запускается после завершения t2 и t3
|
[t2, t3] >> t4 # t4 стартует после завершения и t2, и t3
|
||||||
```
|
```
|
||||||
|
|
||||||
Для сложных пайплайнов можно создавать цепочки любой сложности:
|
Можно комбинировать:
|
||||||
|
|
||||||
```python
|
```python
|
||||||
t1 >> [t2, t3] >> t4 >> [t5, t6, t7] >> t8
|
t1 >> [t2, t3] >> t4 >> [t5, t6, t7] >> t8
|
||||||
```
|
```
|
||||||
|
|
||||||
|
Важно: стрелочный синтаксис хорошо работает для случаев
|
||||||
|
*«одна задача → список задач»* и *«список задач → одна задача»*.
|
||||||
|
|
||||||
|
Но он **не умеет** напрямую связывать два списка между собой:
|
||||||
|
|
||||||
|
```python
|
||||||
|
[t2, t3] >> [t5, t6, t7] # так делать нельзя — будет ошибка
|
||||||
|
```
|
||||||
|
|
||||||
|
Если вам нужно, чтобы **каждая** из задач `t2` и `t3` была предком для **каждой** из задач `t5`, `t6`, `t7`, используйте встроенную функцию `cross_downstream`:
|
||||||
|
|
||||||
|
```python
|
||||||
|
from airflow.models.baseoperator import cross_downstream
|
||||||
|
|
||||||
|
cross_downstream(
|
||||||
|
from_tasks=[t2, t3],
|
||||||
|
to_tasks=[t5, t6, t7],
|
||||||
|
)
|
||||||
|
```
|
||||||
|
|
||||||
|
Такой код создаст зависимости:
|
||||||
|
|
||||||
|
* `t2` → `t5`, `t6`, `t7`
|
||||||
|
* `t3` → `t5`, `t6`, `t7`
|
||||||
|
|
||||||
|
Под капотом `cross_downstream` как раз делает вложенный цикл,
|
||||||
|
но в коде явно видно, что мы хотим «полный крест» между двумя наборами задач, и не приходится писать ручные `for`-ы — это рекомендованный в документации Apache Airflow подход.
|
||||||
|
|
||||||
|
> 💡 В более больших DAG-ах, где таких блоков много, удобнее не оперировать списками, а **группировать задачи в `TaskGroup`** (Task Group). Тогда зависимости задаются уже между группами, а не между отдельными списками задач. Об этом отдельно поговорим в разделе про TaskGroup.
|
||||||
|
|
||||||
## Полный пример DAG-файла
|
## Полный пример DAG-файла
|
||||||
|
|
||||||
В этом полном примере мы создаем DAG с более подробной настройкой параметров. Обратите внимание на дополнительные параметры, такие как количество повторных попыток ('retries'), задержка между попытками ('retry_delay'), а также настройки уведомлений по электронной почте.
|
В этом полном примере мы создаем DAG с более подробной настройкой параметров. Обратите внимание на дополнительные параметры, такие как количество повторных попыток ('retries'), задержка между попытками ('retry_delay'), а также настройки уведомлений по электронной почте.
|
||||||
@@ -178,10 +243,9 @@ t1 >> t2
|
|||||||
2. Создание агрегированной таблицы по регионам и категориям
|
2. Создание агрегированной таблицы по регионам и категориям
|
||||||
3. Сохранение результата в базу данных PostgreSQL
|
3. Сохранение результата в базу данных PostgreSQL
|
||||||
|
|
||||||
Вот полный код нашего DAG:
|
Вот полный код нашего DAG, адаптированный под учебный стенд из папки `airflow-docker`:
|
||||||
|
|
||||||
```python
|
```python
|
||||||
import os
|
|
||||||
import datetime as dt
|
import datetime as dt
|
||||||
import pandas as pd
|
import pandas as pd
|
||||||
from airflow.models import DAG
|
from airflow.models import DAG
|
||||||
@@ -191,6 +255,10 @@ from airflow.operators.bash import BashOperator
|
|||||||
from sqlalchemy import create_engine
|
from sqlalchemy import create_engine
|
||||||
|
|
||||||
|
|
||||||
|
# Подключение к учебной базе PostgreSQL (postgres-training)
|
||||||
|
DB_URL = "postgresql://student:student@postgres-training:5432/training"
|
||||||
|
|
||||||
|
|
||||||
# Базовые параметры DAG
|
# Базовые параметры DAG
|
||||||
args = {
|
args = {
|
||||||
'owner': 'airflow',
|
'owner': 'airflow',
|
||||||
@@ -199,25 +267,30 @@ args = {
|
|||||||
'retry_delay': dt.timedelta(minutes=1),
|
'retry_delay': dt.timedelta(minutes=1),
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
def download_titanic_dataset():
|
def download_titanic_dataset():
|
||||||
|
"""Загрузка датасета Titanic и сохранение в базу"""
|
||||||
url = 'https://web.stanford.edu/class/archive/cs/cs109/cs109.1166/stuff/titanic.csv'
|
url = 'https://web.stanford.edu/class/archive/cs/cs109/cs109.1166/stuff/titanic.csv'
|
||||||
df = pd.read_csv(url)
|
df = pd.read_csv(url)
|
||||||
engine = create_engine('postgresql+psycopg2://jovyan:jovyan@localhost:5432/de')
|
|
||||||
|
engine = create_engine(DB_URL)
|
||||||
df.to_sql('titanic', engine, index=False, if_exists='replace', schema='public')
|
df.to_sql('titanic', engine, index=False, if_exists='replace', schema='public')
|
||||||
|
|
||||||
|
|
||||||
def pivot_dataset():
|
def pivot_dataset():
|
||||||
engine = create_engine('postgresql+psycopg2://jovyan:jovyan@localhost:5432/de')
|
"""Построение сводной таблицы и сохранение результата"""
|
||||||
|
engine = create_engine(DB_URL)
|
||||||
titanic_df = pd.read_sql('select * from public.titanic', con=engine)
|
titanic_df = pd.read_sql('select * from public.titanic', con=engine)
|
||||||
|
|
||||||
df = titanic_df.pivot_table(
|
df = titanic_df.pivot_table(
|
||||||
index=['Sex'],
|
index=['Sex'],
|
||||||
columns=['Pclass'],
|
columns=['Pclass'],
|
||||||
values='Name',
|
values='Name',
|
||||||
aggfunc='count'
|
aggfunc='count'
|
||||||
).reset_index()
|
).reset_index()
|
||||||
|
|
||||||
|
df.to_sql('titanic_pivot', engine, index=False, if_exists='replace', schema='public')
|
||||||
|
|
||||||
df.to_sql('titanic_pivot', engine, index=False, if_exists='replace', schema='public' )
|
|
||||||
|
|
||||||
dag = DAG(
|
dag = DAG(
|
||||||
dag_id='titanic_pivot',
|
dag_id='titanic_pivot',
|
||||||
@@ -228,7 +301,7 @@ dag = DAG(
|
|||||||
# Начальная задача для логирования
|
# Начальная задача для логирования
|
||||||
start = BashOperator(
|
start = BashOperator(
|
||||||
task_id='start',
|
task_id='start',
|
||||||
bash_command='echo "Начинаем выполнение пайплайна! "',
|
bash_command='echo "Начинаем выполнение пайплайна!"',
|
||||||
dag=dag,
|
dag=dag,
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -250,37 +323,23 @@ pivot_titanic_dataset = PythonOperator(
|
|||||||
start >> create_titanic_dataset >> pivot_titanic_dataset
|
start >> create_titanic_dataset >> pivot_titanic_dataset
|
||||||
```
|
```
|
||||||
|
|
||||||
Этот код использует `PythonOperator` для выполнения функций работы с данными, что является стандартной практикой для задач обработки данных. В примере функция `load_customer_dataset` загружает данные из внешнего источника, а `aggregate_customer_dataset` создает агрегированную таблицу.
|
Этот код использует `PythonOperator` для выполнения функций работы с данными: `download_titanic_dataset` загружает исходный датасет и сохраняет его в учебную базу PostgreSQL, а `pivot_dataset` строит сводную таблицу и записывает результат в отдельную таблицу `titanic_pivot`.
|
||||||
|
|
||||||
## Загрузка DAG в учебную среду
|
## Загрузка DAG в учебную среду
|
||||||
|
|
||||||
Для тестирования нашего DAG в учебной среде выполните следующие шаги:
|
Для тестирования нашего DAG в учебной среде выполните следующие шаги:
|
||||||
|
|
||||||
1. Сохраните код в файл с расширением `.py` (например, `customer_analysis_dag.py`)
|
1. Сохраните код в файл с расширением `.py` в папке стенда, например `airflow-docker/dags/titanic_pivot_dag.py`.
|
||||||
|
2. Убедитесь, что стенд запущен:
|
||||||
Для тестирования нашего DAG в учебной среде выполните следующие шаги:
|
|
||||||
|
|
||||||
1. Сохраните код в файл с расширением `.py` (например, `customer_analysis_dag.py`)
|
|
||||||
|
|
||||||
2. Найдите запущенный контейнер с учебной средой Airflow:
|
|
||||||
```bash
|
```bash
|
||||||
docker ps
|
cd airflow-docker
|
||||||
|
docker-compose up -d
|
||||||
```
|
```
|
||||||
|
3. Подождите 30–60 секунд — Airflow автоматически обнаружит новый файл в папке `dags`.
|
||||||
|
4. Откройте веб-интерфейс Airflow (http://localhost:8080), найдите DAG `titanic_pivot` по идентификатору и запустите его.
|
||||||
|
|
||||||
3. Скопируйте файл в контейнер:
|
После запуска вы можете посмотреть статус задач и логи в интерфейсе Airflow
|
||||||
```bash
|
(подробнее про это — в разделе про пользовательский интерфейс).
|
||||||
docker cp /путь/к/файлу/customer_analysis_dag.py [ID_КОНТЕЙНЕРА]:/lessons/dags/customer_analysis_dag.py
|
|
||||||
```
|
|
||||||
|
|
||||||
Например:
|
В этом уроке вы изучили основную структуру DAG-файлов, создали свой первый рабочий пайплайн,
|
||||||
```bash
|
научились загружать его в учебную среду и запускать для получения результатов.
|
||||||
docker cp ~/Desktop/customer_analysis_dag.py 4dc4fbfedcf4:/lessons/dags/customer_analysis_dag.py
|
|
||||||
```
|
|
||||||
|
|
||||||
4. Подождите 30-60 секунд — Airflow автоматически обнаружит новый файл
|
|
||||||
|
|
||||||
5. Найдите ваш DAG в интерфейсе Airflow через строку поиска и запустите его
|
|
||||||
|
|
||||||

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

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

|
На диаграмме ниже показаны те же этапы, но уже с привязкой к внутренним компонентам Airflow.
|
||||||
|
|
||||||
|
```mermaid
|
||||||
|
flowchart LR
|
||||||
|
%% Стили
|
||||||
|
classDef component fill:#8BC34A,stroke:#333,stroke-width:1px,color:#fff;
|
||||||
|
classDef state fill:#ffffff,stroke:#333,stroke-width:1px,color:#000;
|
||||||
|
classDef success fill:#C8E6C9,stroke:#388E3C,stroke-width:1px,color:#000;
|
||||||
|
classDef sensor fill:#ffffff,stroke:#8BC34A,stroke-width:1px,color:#000;
|
||||||
|
|
||||||
|
%% Жизненный цикл задачи
|
||||||
|
|
||||||
|
NONE["No status / None"]:::state --> SCH[Scheduler]:::component
|
||||||
|
|
||||||
|
SCH --> SCHEDULED[Scheduled]:::state
|
||||||
|
SCH --> REMOVED[Removed]:::state
|
||||||
|
SCH --> UPSTREAM_FAILED[Upstream failed]:::state
|
||||||
|
|
||||||
|
SCHEDULED --> EX[Executor]:::component
|
||||||
|
EX --> QUEUED[Queued]:::state
|
||||||
|
QUEUED --> WORKER[Worker]:::component
|
||||||
|
WORKER --> RUNNING[Running]:::state
|
||||||
|
|
||||||
|
RUNNING --> SUCCESS[Success]:::success
|
||||||
|
RUNNING --> FAILED[Failed]:::state
|
||||||
|
RUNNING --> SHUTDOWN[Shutdown]:::state
|
||||||
|
|
||||||
|
%% Альтернативные переходы при ретраях
|
||||||
|
FAILED -.-> UP_FOR_RETRY["Up for retry"]:::state
|
||||||
|
SHUTDOWN -.-> UP_FOR_RETRY
|
||||||
|
UP_FOR_RETRY --> SCH
|
||||||
|
|
||||||
|
%% Сенсоры: режим reschedule
|
||||||
|
RUNNING --> UP_FOR_RESCHEDULE["Up for reschedule"]:::sensor
|
||||||
|
UP_FOR_RESCHEDULE --> SCH
|
||||||
|
|
||||||
|
%% Все стрелки чуть толще
|
||||||
|
linkStyle default stroke-width:2px;
|
||||||
|
|
||||||
|
%% Альтернативные переходы (ретраи) — зелёные и ещё чуть толще
|
||||||
|
linkStyle 11 stroke:#4CAF50,stroke-width:2.5px;
|
||||||
|
linkStyle 12 stroke:#4CAF50,stroke-width:2.5px;
|
||||||
|
```
|
||||||
|
### Легенда к диаграмме состояний задачи
|
||||||
|
|
||||||
|
**Зелёные блоки (Component)** — внутренние компоненты Airflow, которые двигают задачу по жизненному циклу:
|
||||||
|
|
||||||
|
- **Scheduler** — планировщик, проверяет зависимости задач и решает, что ставить в очередь.
|
||||||
|
- **Executor** — исполнитель, получает от Scheduler список задач и распределяет их по воркерам.
|
||||||
|
- **Worker** — рабочий процесс (worker), который фактически запускает код оператора.
|
||||||
|
|
||||||
|
**Белые блоки (Task stage)** — состояния конкретного запуска задачи (task instance, `TaskInstanceState`). На диаграмме они показывают, через какие шаги проходит задача от появления в DAG до успешного завершения или ошибки. Подробный справочник по всем состояниям есть в документации по ссылке ниже.
|
||||||
|
|
||||||
|
**Белые блоки с зелёной рамкой (Task stage only for sensor)** — состояния, характерные только для сенсоров:
|
||||||
|
|
||||||
|
- **Up for reschedule (`up_for_reschedule`)** — сенсор в режиме `reschedule`: условие ещё не выполнено, задача «усыплена» и позже будет снова поставлена в расписание без непрерывной работы воркера.
|
||||||
|
|
||||||
|
> Официальное описание всех состояний `TaskInstanceState` и их жизненного цикла смотрите в документации Airflow:
|
||||||
|
> [Tasks → Task Instances](https://airflow.apache.org/docs/apache-airflow/stable/core-concepts/tasks.html#task-instances).
|
||||||
|
|
||||||
|
|
||||||
## Что делать, если задача упала?
|
## Что делать, если задача упала?
|
||||||
|
|
||||||
@@ -75,4 +132,8 @@
|
|||||||
- **Остальные статусы** вы будете изучать по мере необходимости
|
- **Остальные статусы** вы будете изучать по мере необходимости
|
||||||
- **Цвета в интерфейсе** — это быстрый способ понять состояние вашего пайплайна
|
- **Цвета в интерфейсе** — это быстрый способ понять состояние вашего пайплайна
|
||||||
|
|
||||||
Помните: понимание статусов задач — это как научиться читать дорожные знаки. Сначала кажется много информации, но со временем это становится второй натурой!
|
Помните: понимание статусов задач — это как научиться читать дорожные знаки. Сначала кажется много информации, но со временем это становится второй натурой!
|
||||||
|
|
||||||
|
> Полный список возможных состояний задач (TaskInstanceState) и их классификацию на терминальные и промежуточные можно посмотреть в официальной документации Airflow:
|
||||||
|
> https://airflow.apache.org/docs/apache-airflow/2.9.3/_api/airflow/utils/state/index.html
|
||||||
|
|
||||||
|
|||||||
@@ -1,4 +1,5 @@
|
|||||||
# Гибкие шаблоны и настройки в Airflow
|
# Шаблоны, переменные и подключения в Airflow
|
||||||
|
|
||||||
В этом материале вы познакомитесь с мощными инструментами Airflow для создания гибких и переиспользуемых пайплайнов: динамическими шаблонами, безопасными переменными и централизованными подключениями к внешним системам.
|
В этом материале вы познакомитесь с мощными инструментами Airflow для создания гибких и переиспользуемых пайплайнов: динамическими шаблонами, безопасными переменными и централизованными подключениями к внешним системам.
|
||||||
|
|
||||||
# Динамические шаблоны Airflow (на основе Jinja)
|
# Динамические шаблоны Airflow (на основе Jinja)
|
||||||
@@ -18,30 +19,30 @@
|
|||||||
...
|
...
|
||||||
|
|
||||||
dag = DAG(
|
dag = DAG(
|
||||||
dag_id="dynamic_templates_example",
|
dag_id="template_example",
|
||||||
schedule_interval="*/10 * * * *",
|
schedule_interval="*/15 * * * *",
|
||||||
default_args=default_args
|
default_args=default_args
|
||||||
)
|
)
|
||||||
|
|
||||||
t1 = BashOperator(task_id="show_date", bash_command="echo {{ ds }}")
|
t1 = BashOperator(task_id="display_date", bash_command="echo {{ ds }}")
|
||||||
|
|
||||||
t1
|
t1
|
||||||
```
|
```
|
||||||
|
|
||||||
**Использование шаблонов в идентификаторах задач для лучшей отслеживаемости:**
|
**Использование шаблонов для лучшей отслеживаемости задач:**
|
||||||
|
|
||||||
В этом примере создается DAG с идентификатором "template_tracking_example", который запускается каждые 20 минут. Вторая задача использует шаблон {{ ds }} в своем идентификаторе, что позволяет легко идентифицировать задачу по дате выполнения.
|
В этом примере создается DAG с идентификатором "template_tracking_example", который запускается каждые 20 минут. Вторая задача использует шаблон {{ ds }} в команде bash, чтобы явно указывать дату обработки в логах.
|
||||||
```python
|
```python
|
||||||
...
|
...
|
||||||
|
|
||||||
dag = DAG(
|
dag = DAG(
|
||||||
dag_id="dynamic_templates_example",
|
dag_id="template_tracking_example",
|
||||||
schedule_interval="*/10 * * * *",
|
schedule_interval="*/20 * * * *",
|
||||||
default_args=default_args
|
default_args=default_args
|
||||||
)
|
)
|
||||||
|
|
||||||
t1 = BashOperator(task_id="show_date", bash_command="echo {{ ds }}")
|
t1 = BashOperator(task_id="show_date", bash_command="echo {{ ds }}")
|
||||||
t2 = BashOperator(task_id=f"process_for_{{ ds }}", bash_command="echo {{ ds }}")
|
t2 = BashOperator(task_id="process_for_date", bash_command="echo Processing for {{ ds }}")
|
||||||
|
|
||||||
t1 >> t2
|
t1 >> t2
|
||||||
```
|
```
|
||||||
@@ -50,22 +51,52 @@ t1 >> t2
|
|||||||
|
|
||||||
Переменные Airflow представляют собой пары "ключ-значение", хранящиеся в метадатабазе системы. Они идеально подходят для хранения конфигурационных параметров, таких как пути к скриптам, имена таблиц или другие настройки, которые должны быть доступны в разных DAG.
|
Переменные Airflow представляют собой пары "ключ-значение", хранящиеся в метадатабазе системы. Они идеально подходят для хранения конфигурационных параметров, таких как пути к скриптам, имена таблиц или другие настройки, которые должны быть доступны в разных DAG.
|
||||||
|
|
||||||
Управление переменными осуществляется через веб-интерфейс Airflow (Admin → Variables), где можно:
|
Управление переменными осуществляется через веб-интерфейс Airflow (раздел **Admin → Variables**). Через этот раздел можно:
|
||||||
- Создавать и редактировать пары ключ-значение вручную
|
|
||||||
- Импортировать настройки из JSON-файлов
|
|
||||||
- Использовать командную строку Airflow
|
|
||||||
|
|
||||||

|
- создавать и редактировать пары «ключ-значение» вручную;
|
||||||
|
- импортировать набор переменных из JSON-файла;
|
||||||
|
- удалять больше не нужные настройки.
|
||||||
|
|
||||||
Для защиты конфиденциальной информации Airflow автоматически маскирует значения переменных, в названии которых содержится слово "secret".
|
### Как создать переменную через UI
|
||||||
|
|
||||||
|
Интерфейс ниже соответствует Airflow 2.9.x:
|
||||||
|
|
||||||
|
1. Откройте веб-интерфейс Airflow и авторизуйтесь под пользователем с правами **Admin**.
|
||||||
|
2. В верхнем меню выберите **Admin → Variables**.
|
||||||
|
3. В правом верхнем углу нажмите кнопку **+ Add a new record** (или иконку `+`).
|
||||||
|
4. В поле **Key** задайте имя переменной.
|
||||||
|
Например, создадим переменную с паролем к учебной БД отчётности PostgreSQL:
|
||||||
|
|
||||||
|
- **Key**: `reporting_db_password`
|
||||||
|
5. В поле **Value** введите значение.
|
||||||
|
Например:
|
||||||
|
|
||||||
|
- **Value**: `airflow_report_ro`
|
||||||
|
6. Поле **Description** можно использовать для короткого пояснения, зачем нужна переменная, например:
|
||||||
|
`Пароль read-only к учебной БД отчётности`.
|
||||||
|
7. Нажмите **Save**.
|
||||||
|
|
||||||
|
После сохранения переменная появится в таблице. Значение будет частично скрыто в UI: вместо реального пароля вы увидите `***` — Airflow маскирует секреты в интерфейсе и логах, чтобы их нельзя было случайно подсмотреть.
|
||||||
|
|
||||||
|
Теперь эту переменную можно использовать в коде DAG, например:
|
||||||
|
|
||||||
|
```python
|
||||||
|
from airflow.models import Variable
|
||||||
|
|
||||||
|
reporting_db_password = Variable.get("reporting_db_password")
|
||||||
|
```
|
||||||
|
|
||||||
|
Для защиты конфиденциальной информации Airflow автоматически маскирует значения переменных, в названии которых содержится слово `secret`, а также ряд других чувствительных паттернов.
|
||||||
|
|
||||||
|
Подробнее о переменных — в [официальной документации Airflow](https://airflow.apache.org/docs/apache-airflow/2.9.3/howto/variable.html).
|
||||||
|
|
||||||
**Пример использования переменной в коде DAG:**
|
**Пример использования переменной в коде DAG:**
|
||||||
|
|
||||||
В этом примере создается DAG с идентификатором "variable_example", который использует переменную 'data_storage_path', предварительно сохраненную в Airflow. Значение переменной извлекается с помощью Variable.get() и используется в команде bash для указания пути к данным.
|
В этом примере создается DAG с идентификатором "variable_example", который использует переменную 'data_storage_path', предварительно сохраненную в Airflow. Значение переменной извлекается с помощью Variable.get() и используется в команде bash для указания пути к данным.
|
||||||
|
|
||||||
```python
|
```python
|
||||||
from airflow import DAG
|
from airflow import DAG
|
||||||
from airflow.operators.bash import BashOperator
|
from airflow.operators.bash import BashOperator
|
||||||
from airflow.operators.dummy import DummyOperator
|
|
||||||
from airflow.models import Variable
|
from airflow.models import Variable
|
||||||
from datetime import datetime
|
from datetime import datetime
|
||||||
|
|
||||||
@@ -108,7 +139,7 @@ task = BashOperator(
|
|||||||
В приведенном примере используется PostgresOperator для создания таблицы в базе данных. Вместо использования подключения по умолчанию, явно указывается подключение с идентификатором 'my_postgres_conn'.
|
В приведенном примере используется PostgresOperator для создания таблицы в базе данных. Вместо использования подключения по умолчанию, явно указывается подключение с идентификатором 'my_postgres_conn'.
|
||||||
|
|
||||||
```python
|
```python
|
||||||
from airflow.operators.postgres_operator import PostgresOperator
|
from airflow.providers.postgres.operators.postgres import PostgresOperator
|
||||||
|
|
||||||
create_table = PostgresOperator(
|
create_table = PostgresOperator(
|
||||||
task_id='create_user_table',
|
task_id='create_user_table',
|
||||||
@@ -123,11 +154,53 @@ create_table = PostgresOperator(
|
|||||||
|
|
||||||
Управление подключениями доступно через интерфейс Airflow (Admin → Connections). Если требуемый тип подключения отсутствует, его можно добавить установкой соответствующего Airflow Provider из [официального репозитория](https://airflow.apache.org/docs/#providers-packages-docs-apache-airflow-providers-index-html).
|
Управление подключениями доступно через интерфейс Airflow (Admin → Connections). Если требуемый тип подключения отсутствует, его можно добавить установкой соответствующего Airflow Provider из [официального репозитория](https://airflow.apache.org/docs/#providers-packages-docs-apache-airflow-providers-index-html).
|
||||||
|
|
||||||

|
Управление подключениями доступно через интерфейс Airflow (раздел **Admin → Connections**). Подключения хранятся в метадатабазе Airflow, а пароли и другие чувствительные поля шифруются с помощью Fernet и маскируются в UI и логах.
|
||||||
|
|
||||||
|
### Как создать подключение к PostgreSQL через UI
|
||||||
|
|
||||||
|
Интерфейс ниже соответствует Airflow 2.9.x и стандартному Docker-стенду из документации:
|
||||||
|
|
||||||
|
1. Откройте веб-интерфейс Airflow.
|
||||||
|
2. В верхнем меню выберите **Admin → Connections**.
|
||||||
|
3. В правом верхнем углу нажмите кнопку **+ Add a new record**.
|
||||||
|
4. В форме укажите параметры:
|
||||||
|
|
||||||
|
- **Connection Id**: `my_postgres_conn`
|
||||||
|
Это имя мы будем использовать в коде DAG (параметр `postgres_conn_id`).
|
||||||
|
- **Connection Type**: `Postgres`
|
||||||
|
- **Host**: `postgres`
|
||||||
|
(так называется контейнер PostgreSQL в типовом `docker-compose.yaml` из официальной инструкции).
|
||||||
|
- **Schema**: `airflow`
|
||||||
|
(имя базы данных; в учебном стенде можно использовать стандартную БД).
|
||||||
|
- **Login**: `airflow`
|
||||||
|
- **Password**: `airflow`
|
||||||
|
- **Port**: `5432`
|
||||||
|
|
||||||
|
5. Нажмите кнопку **Test** (если доступна) — Airflow попробует подключиться к базе.
|
||||||
|
6. Если тест успешен, нажмите **Save**.
|
||||||
|
|
||||||
|
Теперь подключение с идентификатором `my_postgres_conn` доступно во всех DAG’ах. В примере ниже PostgresOperator явно использует это подключение:
|
||||||
|
|
||||||
|
```python
|
||||||
|
from airflow.providers.postgres.operators.postgres 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",
|
||||||
|
)
|
||||||
|
```
|
||||||
|
|
||||||
Для успешной работы с внешними системами сначала необходимо создать соответствующее подключение, а затем использовать его идентификатор в операторах вашего DAG.
|
Для успешной работы с внешними системами сначала необходимо создать соответствующее подключение, а затем использовать его идентификатор в операторах вашего DAG.
|
||||||
|
|
||||||
Более подробную информацию о настройке подключений можно найти в [документации Airflow](https://airflow.apache.org/docs/apache-airflow/stable/howto/connection.html).
|
Более подробную информацию о настройке подключений можно найти в
|
||||||
|
[документации Airflow](https://airflow.apache.org/docs/apache-airflow/2.9.3/howto/connection.html).
|
||||||
|
|
||||||
|
|
||||||
# Проверочный список для качественного DAG
|
# Проверочный список для качественного DAG
|
||||||
|
|
||||||
@@ -139,4 +212,4 @@ create_table = PostgresOperator(
|
|||||||
- Требуется ли маскировка конфиденциальных значений?
|
- Требуется ли маскировка конфиденциальных значений?
|
||||||
- Используются ли в логике даты или временные метки, которые можно заменить на шаблоны?
|
- Используются ли в логике даты или временные метки, которые можно заменить на шаблоны?
|
||||||
|
|
||||||
Ответы на эти вопросы помогут вам создавать надежные, безопасные и легко поддерживаемые пайплайны в Airflow.
|
Ответы на эти вопросы помогут вам создавать надежные, безопасные и легко поддерживаемые пайплайны в Airflow.
|
||||||
@@ -4,7 +4,7 @@
|
|||||||
|
|
||||||
# Настройка автоматических запусков
|
# Настройка автоматических запусков
|
||||||
|
|
||||||
Когда вы создаете свой первый пайплайн, вы уже сталкивались с возможностью запускать DAG по расписанию через параметр `schedule_interval`. По умолчанию этот параметр равен `None`, что означает ручной запуск без автоматического расписания.
|
Когда вы создаете свой первый пайплайн, вы уже сталкивались с возможностью запускать DAG по расписанию через параметр `schedule_interval`. По умолчанию этот параметр равен `None`, что означает ручной запуск без автоматического расписания. Начиная с Airflow 2.2 (и, конечно, в 2.9) можно использовать более современный алиас `schedule`, который полностью эквивалентен `schedule_interval`; в этом пособии мы продолжаем использовать `schedule_interval`, чтобы сохранить единый стиль примеров.
|
||||||
|
|
||||||
В приведенном примере создается DAG с идентификатором "daily_data_processing", который будет запускаться ежедневно. DAG использует дату начала 1 января 2020 года и параметр catchup=False, что означает, что пропущенные запуски обрабатываться не будут. Владелец DAG - команда data_team, и для задач в DAG установлена одна попытка повторного запуска при ошибках.
|
В приведенном примере создается DAG с идентификатором "daily_data_processing", который будет запускаться ежедневно. DAG использует дату начала 1 января 2020 года и параметр catchup=False, что означает, что пропущенные запуски обрабатываться не будут. Владелец DAG - команда data_team, и для задач в DAG установлена одна попытка повторного запуска при ошибках.
|
||||||
|
|
||||||
@@ -21,7 +21,7 @@ dag = DAG(
|
|||||||
)
|
)
|
||||||
```
|
```
|
||||||
|
|
||||||
Важно понимать, что Airflow требует указания даты начала работы (`start_date`) для любого DAG с расписанием. Система использует эту дату как отправную точку и планирует первый запуск, добавляя к ней интервал из `schedule_interval`.
|
Важно понимать, что Airflow требует указания даты начала работы (`start_date`) для любого DAG с расписанием. `start_date` задаёт начало первого интервала данных (логическую дату запуска), а фактический старт выполнения происходит после окончания этого интервала. Например, при `@daily` первый запуск с логической датой `2020-01-01` фактически произойдёт около полуночи `2020-01-02` по временной зоне DAG. В продакшен‑практике рекомендуется использовать для `start_date` даты с явной временной зоной (timezone-aware), например на базе библиотеки `pendulum` с таймзоной UTC.
|
||||||
|
|
||||||
## Готовые шаблоны расписания
|
## Готовые шаблоны расписания
|
||||||
|
|
||||||
@@ -44,8 +44,6 @@ Airflow предоставляет удобные встроенные шабл
|
|||||||
|
|
||||||
В этом примере создается DAG с идентификатором "historical_data_processing", который запускается ежедневно и имеет дату начала 1 января 2021 года. Параметр catchup=True означает, что Airflow будет автоматически запускать DAG для всех пропущенных дней с указанной даты начала до текущего момента. Владелец DAG - команда analytics_team, и для задач установлено две попытки повторного запуска при ошибках.
|
В этом примере создается DAG с идентификатором "historical_data_processing", который запускается ежедневно и имеет дату начала 1 января 2021 года. Параметр catchup=True означает, что Airflow будет автоматически запускать DAG для всех пропущенных дней с указанной даты начала до текущего момента. Владелец DAG - команда analytics_team, и для задач установлено две попытки повторного запуска при ошибках.
|
||||||
|
|
||||||
Когда вы создаете DAG с исторической датой начала, Airflow предлагает мощный механизм автоматического пересчета пропущенных периодов через параметр `catchup`.
|
|
||||||
|
|
||||||
```python
|
```python
|
||||||
dag = DAG(
|
dag = DAG(
|
||||||
dag_id="historical_data_processing",
|
dag_id="historical_data_processing",
|
||||||
@@ -61,7 +59,7 @@ dag = DAG(
|
|||||||
|
|
||||||
При `catchup=True` система автоматически выполнит все пропущенные запуски от указанной даты начала до текущего момента. Это особенно полезно при первом запуске DAG для обработки накопившихся исторических данных.
|
При `catchup=True` система автоматически выполнит все пропущенные запуски от указанной даты начала до текущего момента. Это особенно полезно при первом запуске DAG для обработки накопившихся исторических данных.
|
||||||
|
|
||||||
Если ваш бизнес-сценарий не требует пересчета истории или вы хотите начать обработку только с текущего периода, установите `catchup=False`. В этом случае Airflow будет планировать только ближайшие запуски согласно расписанию.
|
По умолчанию для DAG с расписанием параметр `catchup` включен (`True`), поэтому при первом деплое с исторической `start_date` можно неожиданно получить большое количество запусков. Если ваш бизнес-сценарий не требует пересчета истории или вы хотите начать обработку только с текущего периода, установите `catchup=False`. В этом случае Airflow будет планировать только ближайшие запуски согласно расписанию.
|
||||||
|
|
||||||
## Ручная перезаливка данных (backfill)
|
## Ручная перезаливка данных (backfill)
|
||||||
|
|
||||||
@@ -82,7 +80,7 @@ the_main_dag
|
|||||||
|
|
||||||
Вы можете:
|
Вы можете:
|
||||||
- Настроить глобальную временную зону в конфигурационном файле Airflow
|
- Настроить глобальную временную зону в конфигурационном файле Airflow
|
||||||
- Указать временную зону явно при инициализации DAG через параметр `tz`
|
- Указать временную зону для DAG через параметр `timezone` или передать в `start_date` объект с явной таймзоной (например, созданный через `pendulum.datetime(..., tz="UTC")`)
|
||||||
|
|
||||||
Однако рекомендуется придерживаться UTC во всех расчетах и преобразовывать временные метки только при выводе результатов для конечных пользователей. Это минимизирует ошибки и упрощает отладку.
|
Однако рекомендуется придерживаться UTC во всех расчетах и преобразовывать временные метки только при выводе результатов для конечных пользователей. Это минимизирует ошибки и упрощает отладку.
|
||||||
|
|
||||||
@@ -93,4 +91,4 @@ the_main_dag
|
|||||||
3. **Планируйте запуски с учетом UTC**, особенно если ваша команда работает в разных часовых поясах
|
3. **Планируйте запуски с учетом UTC**, особенно если ваша команда работает в разных часовых поясах
|
||||||
4. **Документируйте временные зависимости** в коде DAG для других разработчиков
|
4. **Документируйте временные зависимости** в коде DAG для других разработчиков
|
||||||
|
|
||||||
Понимание временных механизмов Airflow — ключ к созданию надежных и предсказуемых пайплайнов. Правильная настройка расписаний и управление историческими данными позволяют автоматизировать сложные бизнес-процессы без ручного вмешательства.
|
Понимание временных механизмов Airflow — ключ к созданию надежных и предсказуемых пайплайнов. Правильная настройка расписаний и управление историческими данными позволяют автоматизировать сложные бизнес-процессы без ручного вмешательства.
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
# Продвинутые возможности Airflow для начинающих специалистов
|
# Пулы, XCom, TaskGroup и алертинг в Airflow
|
||||||
|
|
||||||
В этом материале мы познакомимся с расширенными функциями Airflow, которые помогут вам решать нетривиальные задачи при работе со сложными пайплайнами.
|
В этом материале мы познакомимся с управлением ресурсами через пулы задач, обменом данными между задачами через XCom, группировкой задач с помощью TaskGroup и настройкой системы оповещений (алертинга) в Airflow.
|
||||||
|
|
||||||
Начнем с улучшения нашего базового пайплайна — проведем рефакторинг кода для лучшей читаемости и поддержки.
|
Начнем с улучшения нашего базового пайплайна — проведем рефакторинг кода для лучшей читаемости и поддержки.
|
||||||
|
|
||||||
@@ -12,11 +12,18 @@
|
|||||||
|
|
||||||
Когда планировщик обнаруживает, что наступило время выполнения DAG, он запускает задачу согласно заданной последовательности. При этом задача занимает один слот в пуле и освобождает его после завершения.
|
Когда планировщик обнаруживает, что наступило время выполнения DAG, он запускает задачу согласно заданной последовательности. При этом задача занимает один слот в пуле и освобождает его после завершения.
|
||||||
|
|
||||||
Создать новый пул задач и установить его размер можно через веб-интерфейс Airflow:
|
Создать новый пул задач и установить его размер можно через веб-интерфейс Airflow. Рассмотрим пример пула `data_processing_pool` для тяжёлых задач:
|
||||||
|
|
||||||

|
1. В верхнем меню откройте **Admin → Pools**.
|
||||||
|
2. Нажмите кнопку **+** или **Create**.
|
||||||
|
3. В поле **Name** укажите, например, `data_processing_pool`.
|
||||||
|
4. В поле **Slots** задайте количество слотов, например `5`
|
||||||
|
(это значит, что одновременно смогут выполняться не более 5 задач из этого пула).
|
||||||
|
5. При желании заполните **Description** — например, `Пул для тяжёлых задач бэкапа`.
|
||||||
|
6. Нажмите **Save**.
|
||||||
|
|
||||||
|
После этого в задачах можно указать параметр `pool="data_processing_pool"`, и они будут занимать слоты именно этого пула.
|
||||||
|
|
||||||
*Создание нового пула задач в Airflow*
|
|
||||||
|
|
||||||
Зачем нужны пулы задач? Они помогают:
|
Зачем нужны пулы задач? Они помогают:
|
||||||
- Организовать запуск процессов в системе
|
- Организовать запуск процессов в системе
|
||||||
@@ -59,13 +66,34 @@ BashOperator(
|
|||||||
|
|
||||||
💡 **Контекст выполнения** — это набор параметров, передаваемых при запуске задачи, а также метаданные, генерируемые во время выполнения: время старта, время завершения, имя DAG и другие. Контекст можно найти следующим образом:
|
💡 **Контекст выполнения** — это набор параметров, передаваемых при запуске задачи, а также метаданные, генерируемые во время выполнения: время старта, время завершения, имя DAG и другие. Контекст можно найти следующим образом:
|
||||||
|
|
||||||
{:height 437, :width 778}
|
1. Откройте веб-интерфейс Airflow и перейдите на вкладку **DAGs**.
|
||||||
|
2. Найдите нужный DAG и кликните по его имени.
|
||||||
|
3. На странице DAG убедитесь, что открыта вкладка **Grid**.
|
||||||
|
4. В сетке выберите нужный запуск DAG (колонка) и задачу (строка) и кликните по цветному квадратику задачи.
|
||||||
|
5. Справа откроется панель **Task Instance** (детали экземпляра задачи). В ней можно увидеть:
|
||||||
|
- `dag_id` и `task_id`;
|
||||||
|
- логическую дату запуска (**logical_date / execution_date**);
|
||||||
|
- текущий статус, время старта и завершения;
|
||||||
|
- ссылки на лог, XCom и другую служебную информацию.
|
||||||
|
|
||||||
В XCom можно передавать сериализованные объекты. Значения XCom хранятся в базе данных Airflow и доступны через интерфейс:
|
Именно эти поля и составляют большую часть «контекста выполнения», который доступен в Jinja-шаблонах через объекты вроде `{{ dag_run }}` и `{{ task_instance }}`.
|
||||||
|
|
||||||

|
В XCom можно передавать сериализованные объекты. Значения XCom хранятся в базе данных Airflow и доступны через интерфейс.
|
||||||
|
|
||||||
*XCom в веб-интерфейсе Airflow*
|
Посмотреть XCom можно двумя способами.
|
||||||
|
|
||||||
|
**1. Через конкретную задачу в DAG**
|
||||||
|
|
||||||
|
1. Откройте нужный DAG и вкладку **Grid**.
|
||||||
|
2. Найдите нужный запуск и кликните по квадратику задачи.
|
||||||
|
3. В правой панели **Task Instance** перейдите на вкладку **XCom**.
|
||||||
|
4. В таблице вы увидите все XCom-записи для этого экземпляра задачи: ключ (`key`), значение (`value`), время создания и т. д.
|
||||||
|
|
||||||
|
**2. Через общий список XCom**
|
||||||
|
|
||||||
|
1. В верхнем меню выберите **Browse → XComs**.
|
||||||
|
2. Отфильтруйте записи по `dag_id`, `task_id` или другим полям, если нужно.
|
||||||
|
3. Откройте интересующую запись, чтобы увидеть её содержимое.
|
||||||
|
|
||||||
**Важное правило**: XCom предназначен для обмена небольшими сообщениями. Данные проходят сериализацию/десериализацию при чтении и записи в таблицу.
|
**Важное правило**: XCom предназначен для обмена небольшими сообщениями. Данные проходят сериализацию/десериализацию при чтении и записи в таблицу.
|
||||||
|
|
||||||
@@ -73,7 +101,7 @@ BashOperator(
|
|||||||
|
|
||||||
Для передачи больших объемов данных используйте внешние средства: файловую систему, базы данных (чаще всего PostgreSQL) или очереди сообщений (например, Kafka).
|
Для передачи больших объемов данных используйте внешние средства: файловую систему, базы данных (чаще всего PostgreSQL) или очереди сообщений (например, Kafka).
|
||||||
|
|
||||||
Механизм XCom похож на работу функций в Python. Многие операторы по умолчанию возвращают результат выполнения задачи. За это отвечает параметр `do_xcom_push`, который по умолчанию равен `True`.
|
Механизм XCom похож на работу функций в Python. Многие операторы (например, `PythonOperator`) по умолчанию возвращают результат выполнения задачи. За это отвечает параметр `do_xcom_push`, который во многих случаях равен `True` по умолчанию.
|
||||||
|
|
||||||
Чтобы прочитать сообщения из XCom, используйте метод `xcom_pull` в контексте задачи:
|
Чтобы прочитать сообщения из XCom, используйте метод `xcom_pull` в контексте задачи:
|
||||||
|
|
||||||
@@ -149,17 +177,44 @@ with TaskGroup("data_loading") as loading_group:
|
|||||||
|
|
||||||
Обратите внимание: при использовании TaskGroup последовательность задач указывается внутри группы после объявления всех задач, а в конце DAG описывается последовательность выполнения самих групп.
|
Обратите внимание: при использовании TaskGroup последовательность задач указывается внутри группы после объявления всех задач, а в конце DAG описывается последовательность выполнения самих групп.
|
||||||
|
|
||||||
Визуально в интерфейсе Airflow это выглядит так:
|
Ниже приведена упрощённая схема зависимостей между тремя группами задач:
|
||||||
|
|
||||||

|
```mermaid
|
||||||
|
flowchart LR
|
||||||
|
|
||||||
*Группировка задач с помощью TaskGroup*
|
%% group1
|
||||||
|
subgraph G1["group1"]
|
||||||
|
g1_t1["task1"]
|
||||||
|
g1_t2["task2"]
|
||||||
|
g1_t3["task3"]
|
||||||
|
|
||||||
TaskGroup добавляет интерактивность в веб-интерфейс — группы задач можно сворачивать и разворачивать, что значительно улучшает восприятие DAG с большим количеством задач и связей:
|
g1_t1 --> g1_t2
|
||||||
|
g1_t1 --> g1_t3
|
||||||
|
end
|
||||||
|
|
||||||

|
%% group2
|
||||||
|
subgraph G2["group2"]
|
||||||
|
g2_t1["task1"]
|
||||||
|
g2_t2["task2"]
|
||||||
|
|
||||||
*Группу задач можно раскрыть для детального просмотра*
|
g2_t1 --> g2_t2
|
||||||
|
end
|
||||||
|
|
||||||
|
%% group3
|
||||||
|
subgraph G3["group3"]
|
||||||
|
g3_t1["task1"]
|
||||||
|
g3_t2["task2"]
|
||||||
|
|
||||||
|
g3_t1 --> g3_t2
|
||||||
|
end
|
||||||
|
|
||||||
|
%% зависимости между группами
|
||||||
|
g1_t2 --> g3_t1
|
||||||
|
g1_t3 --> g3_t1
|
||||||
|
g2_t2 --> g3_t1
|
||||||
|
```
|
||||||
|
|
||||||
|
Визуально в интерфейсе Airflow группы задач отображаются как один узел с небольшим индикатором. Клик по нему разворачивает или сворачивает вложенные задачи, что значительно улучшает восприятие DAG с большим количеством задач и связей, особенно когда в них десятки и сотни задач.
|
||||||
|
|
||||||
TaskGroup — это удобный способ логической группировки задач, который помогает упростить код и представить сложные пайплайны более компактно.
|
TaskGroup — это удобный способ логической группировки задач, который помогает упростить код и представить сложные пайплайны более компактно.
|
||||||
|
|
||||||
@@ -167,7 +222,7 @@ TaskGroup — это удобный способ логической групп
|
|||||||
|
|
||||||
Алертинг — один из ключевых компонентов системы оркестрации, так как важно своевременно получать уведомления об ошибках для их оперативного анализа и решения.
|
Алертинг — один из ключевых компонентов системы оркестрации, так как важно своевременно получать уведомления об ошибках для их оперативного анализа и решения.
|
||||||
|
|
||||||
По умолчанию в Airflow настроена отправка уведомлений на электронную почту. При создании DAG указываются email-адреса, на которые будут отправляться сообщения. С помощью параметров можно настроить различные сценарии оповещений.
|
В Airflow есть встроенная поддержка отправки уведомлений на электронную почту (при условии, что в конфигурации настроен SMTP-сервер). При создании DAG указываются email-адреса, на которые будут отправляться сообщения. С помощью параметров можно настроить различные сценарии оповещений.
|
||||||
|
|
||||||
Давайте модифицируем наш первый DAG так, чтобы получать уведомления на почту при возникновении ошибок. При этом настроим перезапуск задач в случае неудачи (например, 2 попытки), но без уведомлений о самих перезапусках:
|
Давайте модифицируем наш первый DAG так, чтобы получать уведомления на почту при возникновении ошибок. При этом настроим перезапуск задач в случае неудачи (например, 2 попытки), но без уведомлений о самих перезапусках:
|
||||||
|
|
||||||
@@ -208,14 +263,16 @@ dag = DAG(
|
|||||||
from datetime import datetime, timedelta, timezone
|
from datetime import datetime, timedelta, timezone
|
||||||
import dateutil
|
import dateutil
|
||||||
from airflow.utils.email import send_email_smtp
|
from airflow.utils.email import send_email_smtp
|
||||||
|
from airflow.operators.python import get_current_context
|
||||||
|
|
||||||
MAIL_LIST = [
|
MAIL_LIST = [
|
||||||
"email_1@gmail.ru",
|
"email_1@gmail.ru",
|
||||||
"email_2@gmail.ru"
|
"email_2@gmail.ru"
|
||||||
]
|
]
|
||||||
|
|
||||||
def notify_email(calculation_dt: str, dagrun_begin_time, **context):
|
def notify_email(calculation_dt: str, dagrun_begin_time: str):
|
||||||
# dag_run date is utc timezone, so add `timezone.utc` to calculation_end to combat 3 hour diff.
|
# dag_run date is utc timezone, so add `timezone.utc` to calculation_end to combat 3 hour diff.
|
||||||
|
context = get_current_context()
|
||||||
calculation_start = dateutil.parser.isoparse(dagrun_begin_time.replace('Z', '+00:00'))
|
calculation_start = dateutil.parser.isoparse(dagrun_begin_time.replace('Z', '+00:00'))
|
||||||
calculation_end = datetime.now(timezone.utc)
|
calculation_end = datetime.now(timezone.utc)
|
||||||
duration = str(timedelta(seconds=(calculation_end - calculation_start).seconds))
|
duration = str(timedelta(seconds=(calculation_end - calculation_start).seconds))
|
||||||
@@ -241,6 +298,7 @@ def notify_email(calculation_dt: str, dagrun_begin_time, **context):
|
|||||||
|
|
||||||
```python
|
```python
|
||||||
from airflow.operators.python import PythonOperator
|
from airflow.operators.python import PythonOperator
|
||||||
|
from airflow.utils.trigger_rule import TriggerRule
|
||||||
from utils import notify_email
|
from utils import notify_email
|
||||||
|
|
||||||
LOCAL_CALCULATION_DT = '{{ dag.timezone.convert(execution_date).strftime("%Y-%m-%d") }}'
|
LOCAL_CALCULATION_DT = '{{ dag.timezone.convert(execution_date).strftime("%Y-%m-%d") }}'
|
||||||
@@ -249,7 +307,6 @@ DAG_RUN_BEGIN_TIME = "{{ dag_run.start_date }}"
|
|||||||
email_notification_task = PythonOperator(
|
email_notification_task = PythonOperator(
|
||||||
task_id="send_email_notification",
|
task_id="send_email_notification",
|
||||||
python_callable=notify_email,
|
python_callable=notify_email,
|
||||||
provide_context=True,
|
|
||||||
dag=dag,
|
dag=dag,
|
||||||
trigger_rule=TriggerRule.ALL_DONE,
|
trigger_rule=TriggerRule.ALL_DONE,
|
||||||
op_args=[LOCAL_CALCULATION_DT, DAG_RUN_BEGIN_TIME],
|
op_args=[LOCAL_CALCULATION_DT, DAG_RUN_BEGIN_TIME],
|
||||||
@@ -267,6 +324,7 @@ email_notification_task = PythonOperator(
|
|||||||
В этом примере используется переменная 'database_connection_string', предварительно созданная в интерфейсе Airflow, для хранения строки подключения к базе данных. Это позволяет избежать жесткого кодирования конфиденциальной информации в коде DAG и упрощает настройку подключения для разных окружений.
|
В этом примере используется переменная 'database_connection_string', предварительно созданная в интерфейсе Airflow, для хранения строки подключения к базе данных. Это позволяет избежать жесткого кодирования конфиденциальной информации в коде DAG и упрощает настройку подключения для разных окружений.
|
||||||
|
|
||||||
```python
|
```python
|
||||||
|
import os
|
||||||
import datetime as dt
|
import datetime as dt
|
||||||
import pandas as pd
|
import pandas as pd
|
||||||
from airflow.models import DAG
|
from airflow.models import DAG
|
||||||
@@ -291,8 +349,8 @@ def get_file_path(file_name):
|
|||||||
return os.path.join(os.path.expanduser('~/data'), file_name)
|
return os.path.join(os.path.expanduser('~/data'), file_name)
|
||||||
|
|
||||||
def load_customer_data():
|
def load_customer_data():
|
||||||
url = 'https://example.com/customer_data.csv'
|
file_path = get_file_path('customer_data.csv')
|
||||||
df = pd.read_csv(url)
|
df = pd.read_csv(file_path)
|
||||||
engine = create_engine(DATABASE_URL)
|
engine = create_engine(DATABASE_URL)
|
||||||
df.to_sql('customers', engine, index=False, if_exists='replace', schema='staging')
|
df.to_sql('customers', engine, index=False, if_exists='replace', schema='staging')
|
||||||
|
|
||||||
@@ -351,6 +409,7 @@ start_task >> data_processing
|
|||||||
Итоговый DAG будет выглядеть следующим образом:
|
Итоговый DAG будет выглядеть следующим образом:
|
||||||
|
|
||||||
```python
|
```python
|
||||||
|
import os
|
||||||
import datetime as dt
|
import datetime as dt
|
||||||
import pandas as pd
|
import pandas as pd
|
||||||
from airflow.models import DAG
|
from airflow.models import DAG
|
||||||
@@ -374,8 +433,8 @@ def get_file_path(file_name):
|
|||||||
return os.path.join(os.path.expanduser('~/data'), file_name)
|
return os.path.join(os.path.expanduser('~/data'), file_name)
|
||||||
|
|
||||||
def load_customer_data():
|
def load_customer_data():
|
||||||
url = 'https://example.com/customer_data.csv'
|
file_path = get_file_path('customer_data.csv')
|
||||||
df = pd.read_csv(url)
|
df = pd.read_csv(file_path)
|
||||||
engine = create_engine(DATABASE_URL)
|
engine = create_engine(DATABASE_URL)
|
||||||
df.to_sql('customers', engine, index=False, if_exists='replace', schema='staging')
|
df.to_sql('customers', engine, index=False, if_exists='replace', schema='staging')
|
||||||
|
|
||||||
@@ -426,4 +485,4 @@ start_task >> data_processing
|
|||||||
- Логическая группировка задач с TaskGroup
|
- Логическая группировка задач с TaskGroup
|
||||||
- Настройка системы оповещений
|
- Настройка системы оповещений
|
||||||
|
|
||||||
Помните: не стоит использовать все доступные функции сразу. Выбирайте инструменты последовательно и находите оптимальный набор возможностей под конкретную задачу.
|
Помните: не стоит использовать все доступные функции сразу. Выбирайте инструменты последовательно и находите оптимальный набор возможностей под конкретную задачу.
|
||||||
@@ -1,6 +1,6 @@
|
|||||||
# Repository Guidelines
|
# Repository Guidelines
|
||||||
|
|
||||||
This repository contains an educational manual for Apache Airflow with runnable examples. The root holds the written guides; `airflow-docker/` provides a self-contained Docker setup to run and explore DAGs.
|
This repository contains an educational manual for Apache Airflow with runnable examples. The root holds the written guides; `airflow-docker/` provides a self-contained Docker setup to run and explore DAGs. The course and examples target Apache Airflow **2.9.x**.
|
||||||
|
|
||||||
## Project Structure & Module Organization
|
## Project Structure & Module Organization
|
||||||
- Root `01-09 *.md`: step-by-step articles (RU).
|
- Root `01-09 *.md`: step-by-step articles (RU).
|
||||||
@@ -41,3 +41,6 @@ Prerequisite: Docker + Docker Compose.
|
|||||||
## Security & Configuration Tips
|
## Security & Configuration Tips
|
||||||
- Do not commit secrets. Use environment variables and local `.env` files if needed.
|
- Do not commit secrets. Use environment variables and local `.env` files if needed.
|
||||||
- Use `airflow-docker/data/` for sample data; avoid real PII in the repo.
|
- Use `airflow-docker/data/` for sample data; avoid real PII in the repo.
|
||||||
|
|
||||||
|
## Hints
|
||||||
|
- При работе под Windows для работы с командной строкой используй PowerShell
|
||||||
|
|||||||
@@ -11,7 +11,7 @@ Apache Airflow — это мощный оркестратор рабочих п
|
|||||||
- Хотите освоить свой первый ETL-инструмент
|
- Хотите освоить свой первый ETL-инструмент
|
||||||
- Нуждаетесь в практическом руководстве с понятными объяснениями
|
- Нуждаетесь в практическом руководстве с понятными объяснениями
|
||||||
|
|
||||||
Мы начнем с основ и постепенно перейдем к продвинутым возможностям, чтобы вы могли уверенно использовать Airflow в реальных проектах. Курс ориентирован на версию Airflow 2.5 и использует практический подход с множеством примеров и визуальных материалов.
|
Мы начнем с основ и постепенно перейдем к продвинутым возможностям, чтобы вы могли уверенно использовать Airflow в реальных проектах. Курс ориентирован на версию Airflow 2.9 и использует практический подход с множеством примеров и визуальных материалов.
|
||||||
|
|
||||||
## 🎯 Что вы узнаете
|
## 🎯 Что вы узнаете
|
||||||
|
|
||||||
@@ -21,7 +21,7 @@ Apache Airflow — это мощный оркестратор рабочих п
|
|||||||
- **Интерфейс**: как эффективно использовать веб-интерфейс для мониторинга и управления
|
- **Интерфейс**: как эффективно использовать веб-интерфейс для мониторинга и управления
|
||||||
- **Практическое создание DAG**: пошаговое руководство по созданию ваших первых рабочих процессов
|
- **Практическое создание DAG**: пошаговое руководство по созданию ваших первых рабочих процессов
|
||||||
- **Отладка и мониторинг**: как понимать статусы задач и быстро находить проблемы
|
- **Отладка и мониторинг**: как понимать статусы задач и быстро находить проблемы
|
||||||
- **Продвинутые возможности**: шаблоны, параметры, управление временем и другие мощные функции
|
- **Продвинутые возможности**: шаблоны, параметры, управление временем, пулы, XCom, TaskGroup и алертинг
|
||||||
|
|
||||||
## 📖 Содержание курса
|
## 📖 Содержание курса
|
||||||
|
|
||||||
@@ -57,20 +57,21 @@ Apache Airflow — это мощный оркестратор рабочих п
|
|||||||
- Жизненный цикл задачи
|
- Жизненный цикл задачи
|
||||||
- Как отлаживать проблемы и работать с ошибками
|
- Как отлаживать проблемы и работать с ошибками
|
||||||
|
|
||||||
### [07. Гибкие шаблоны и настройки в Airflow](07%20-%20Гибкие%20шаблоны%20и%20настройки%20в%20Airflow.md)
|
### [07. Шаблоны, переменные и подключения в Airflow](07%20-%20Шаблоны,%20переменные%20и%20подключения%20в%20Airflow.md)
|
||||||
- Использование Jinja-шаблонов
|
- Использование Jinja-шаблонов
|
||||||
- Параметризация DAG
|
- Параметризация DAG
|
||||||
- Глобальные переменные и соединения
|
- Глобальные переменные и подключения
|
||||||
|
|
||||||
### [08. Управление временем в Airflow](08%20-%20Управление%20временем.md)
|
### [08. Управление временем в Airflow](08%20-%20Управление%20временем.md)
|
||||||
- Расписания и интервалы запуска
|
- Расписания и интервалы запуска
|
||||||
- Работа с временными зонами
|
- Работа с временными зонами
|
||||||
- Понимание execution_date и других временных концепций
|
- Понимание execution_date и других временных концепций
|
||||||
|
|
||||||
### [09. Продвинутые возможности Airflow](09%20-%20Продвинутые%20возможности%20Airflow.md)
|
### [09. Пулы, XCom, TaskGroup и алертинг в Airflow](09%20-%20Пулы,%20XCom,%20TaskGroup%20и%20алертинг%20в%20Airflow.md)
|
||||||
- Расширенные операторы и сенсоры
|
- Управление ресурсами с помощью пулов задач
|
||||||
- Обработка ошибок и повторные попытки
|
- Обмен данными между задачами через XCom
|
||||||
- Оптимизация и масштабирование DAG
|
- Логическая группировка задач с TaskGroup
|
||||||
|
- Настройка системы оповещений и уведомлений
|
||||||
|
|
||||||
## 🛠️ Практический стенд
|
## 🛠️ Практический стенд
|
||||||
|
|
||||||
@@ -102,4 +103,4 @@ Apache Airflow — это мощный оркестратор рабочих п
|
|||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
*Этот учебник создан для начинающих специалистов в области данных и инженерии. Все материалы ориентированы на практическое применение и пошаговое освоение Apache Airflow.*
|
*Этот учебник создан для начинающих специалистов в области данных и инженерии. Все материалы ориентированы на практическое применение и пошаговое освоение Apache Airflow.*
|
||||||
|
|||||||
|
Before Width: | Height: | Size: 52 KiB |
|
Before Width: | Height: | Size: 1.3 MiB |
|
Before Width: | Height: | Size: 1.4 MiB |
|
Before Width: | Height: | Size: 9.6 MiB |
|
Before Width: | Height: | Size: 89 KiB |
|
Before Width: | Height: | Size: 160 KiB |
|
Before Width: | Height: | Size: 25 KiB |
|
Before Width: | Height: | Size: 60 KiB |
|
Before Width: | Height: | Size: 28 KiB |
|
Before Width: | Height: | Size: 778 KiB |
|
Before Width: | Height: | Size: 1.0 MiB |
|
Before Width: | Height: | Size: 34 KiB |
|
Before Width: | Height: | Size: 475 KiB |
|
Before Width: | Height: | Size: 187 KiB |
|
Before Width: | Height: | Size: 311 KiB |
|
Before Width: | Height: | Size: 807 KiB |
|
Before Width: | Height: | Size: 838 KiB |
|
Before Width: | Height: | Size: 880 KiB |
|
Before Width: | Height: | Size: 873 KiB |
|
Before Width: | Height: | Size: 322 KiB |
|
Before Width: | Height: | Size: 834 KiB |
|
Before Width: | Height: | Size: 142 KiB |
|
Before Width: | Height: | Size: 188 KiB |
|
Before Width: | Height: | Size: 649 KiB |
|
Before Width: | Height: | Size: 2.0 MiB |
|
Before Width: | Height: | Size: 502 KiB |
|
Before Width: | Height: | Size: 164 KiB |
|
Before Width: | Height: | Size: 22 KiB |
|
Before Width: | Height: | Size: 546 KiB |
|
Before Width: | Height: | Size: 748 KiB |
|
Before Width: | Height: | Size: 2.4 MiB |
@@ -7,9 +7,8 @@ This file provides guidance for AI agents working on the educational Airflow set
|
|||||||
|
|
||||||
```
|
```
|
||||||
airflow-docker/
|
airflow-docker/
|
||||||
├── docker-compose.yml # Docker configuration
|
├── docker-compose.yml # Docker configuration and environment
|
||||||
├── .env # Hardcoded environment variables
|
├── dags/ # DAG files for learning
|
||||||
├── dags/ # DAG files for learning
|
|
||||||
│ ├── hello_world_dag.py # Basic Python operators
|
│ ├── hello_world_dag.py # Basic Python operators
|
||||||
│ ├── sql_basic_dag.py # SQL operations
|
│ ├── sql_basic_dag.py # SQL operations
|
||||||
│ ├── file_operations_dag.py # File processing
|
│ ├── file_operations_dag.py # File processing
|
||||||
@@ -118,8 +117,7 @@ docker-compose logs postgres-metadata
|
|||||||
## 📚 Key Files to Examine
|
## 📚 Key Files to Examine
|
||||||
|
|
||||||
### Configuration Files
|
### Configuration Files
|
||||||
- [`docker-compose.yml`](docker-compose.yml) - Main Docker configuration
|
- [`docker-compose.yml`](docker-compose.yml) - Main Docker configuration and environment variables
|
||||||
- [`.env`](.env) - Environment variables
|
|
||||||
- [`requirements.txt`](requirements.txt) - Python dependencies
|
- [`requirements.txt`](requirements.txt) - Python dependencies
|
||||||
|
|
||||||
### Sample Data Files
|
### Sample Data Files
|
||||||
@@ -184,4 +182,4 @@ docker-compose logs postgres-metadata
|
|||||||
### For New DAGs
|
### For New DAGs
|
||||||
- Include comprehensive docstring in Russian
|
- Include comprehensive docstring in Russian
|
||||||
- Describe learning objectives clearly
|
- Describe learning objectives clearly
|
||||||
- Provide step-by-step task explanations
|
- Provide step-by-step task explanations
|
||||||
|
|||||||
@@ -131,6 +131,27 @@ Process data based on file type or data quality
|
|||||||
- `success_handler`: On success callback
|
- `success_handler`: On success callback
|
||||||
- `failure_handler`: On failure callback
|
- `failure_handler`: On failure callback
|
||||||
|
|
||||||
|
### Level 4: Orchestration & Collaboration (Week 4)
|
||||||
|
|
||||||
|
#### 4.1 advanced_features_dag.py
|
||||||
|
**Learning Objectives:**
|
||||||
|
- Control concurrency with pools and `pool_slots`
|
||||||
|
- Exchange data between tasks using XCom
|
||||||
|
- Group related tasks using `TaskGroup`
|
||||||
|
- Configure email-based alerting on failures
|
||||||
|
|
||||||
|
**Scenario:**
|
||||||
|
Enhanced daily analytics pipeline that reads data from the training database, performs transformations, and writes summaries, while limiting heavy backup tasks via a dedicated pool and sending notifications about pipeline status.
|
||||||
|
|
||||||
|
**Tasks:**
|
||||||
|
- `extract_group`: Use `TaskGroup` to wrap extract tasks (e.g., customers and orders)
|
||||||
|
- `transform_group`: Aggregate metrics and prepare summary tables
|
||||||
|
- `load_group`: Simulate loading results back into the training database or files
|
||||||
|
- `backup_task`: Heavy backup task running in a dedicated pool (e.g., `backup_pool`) with custom `pool_slots`
|
||||||
|
- `calculate_metrics`: Python task that returns aggregated metrics (pushed to XCom)
|
||||||
|
- `log_metrics`: Task that reads metrics via `xcom_pull` and logs them or uses them in a template
|
||||||
|
- `send_notification`: Final notification task (email or log) triggered with `ALL_DONE` semantics
|
||||||
|
|
||||||
## Sample Data Files
|
## Sample Data Files
|
||||||
|
|
||||||
### customers.csv
|
### customers.csv
|
||||||
@@ -214,6 +235,12 @@ CREATE TABLE enrollments (
|
|||||||
- ✅ Use parameters and templates
|
- ✅ Use parameters and templates
|
||||||
- ✅ Monitor and debug workflows
|
- ✅ Monitor and debug workflows
|
||||||
|
|
||||||
|
### Week 4: Orchestration & Operations
|
||||||
|
- ✅ Use pools to control resource usage
|
||||||
|
- ✅ Share data between tasks via XCom
|
||||||
|
- ✅ Group tasks using `TaskGroup`
|
||||||
|
- ✅ Configure alerting and notifications for failures
|
||||||
|
|
||||||
## Common Pitfalls and Solutions
|
## Common Pitfalls and Solutions
|
||||||
|
|
||||||
### Problem: DAG not appearing in UI
|
### Problem: DAG not appearing in UI
|
||||||
|
|||||||
@@ -245,6 +245,47 @@
|
|||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
|
## 🌟 Дополнительные задания: пулы, XCom, TaskGroup и алертинг
|
||||||
|
|
||||||
|
**Цель:** Освоить продвинутые возможности оркестрации — управление ресурсами, обмен данными между задачами, группировку и оповещения.
|
||||||
|
|
||||||
|
### Задание 1: Пулы и управление ресурсами
|
||||||
|
**Сложность:** 🟡 Средняя
|
||||||
|
**Время выполнения:** 20-25 минут
|
||||||
|
|
||||||
|
**Задача:**
|
||||||
|
- В интерфейсе Airflow создайте пул `backup_pool` с 2 слотами
|
||||||
|
- Создайте новый DAG `advanced_features_dag.py` **или** расширьте `data_processing_dag.py` задачами резервного копирования (например, `backup_to_csv`, `backup_to_db`)
|
||||||
|
- Назначьте этим задачам параметр `pool="backup_pool"` и настройте `pool_slots` так, чтобы одна из задач занимала 2 слота, а другая — 1
|
||||||
|
- Наблюдайте в UI, что одновременно запускается не более 2 задач из этого пула
|
||||||
|
|
||||||
|
**Цель задания:** Научиться управлять параллелизмом задач через пулы и `pool_slots`.
|
||||||
|
|
||||||
|
### Задание 2: Обмен данными через XCom
|
||||||
|
**Сложность:** 🟡 Средняя
|
||||||
|
**Время выполнения:** 20-25 минут
|
||||||
|
|
||||||
|
**Задача:**
|
||||||
|
- В том же DAG добавьте задачу `calculate_metrics` (PythonOperator), которая возвращает словарь с агрегированными показателями, например: `{"total_orders": ..., "avg_amount": ...}`
|
||||||
|
- Добавьте задачу `log_metrics`, которая с помощью `xcom_pull` читает результат `calculate_metrics` и выводит значения в лог
|
||||||
|
- Для одной из задач продемонстрируйте использование XCom в Jinja-шаблоне (например, в `bash_command` или SQL-запросе)
|
||||||
|
|
||||||
|
**Цель задания:** Освоить передачу результатов между задачами через XCom и их использование в шаблонах.
|
||||||
|
|
||||||
|
### Задание 3: TaskGroup и алертинг
|
||||||
|
**Сложность:** 🟠 Продвинутая
|
||||||
|
**Время выполнения:** 30-35 минут
|
||||||
|
|
||||||
|
**Задача:**
|
||||||
|
- Объедините логически связанные задачи (например, `extract` / `transform` / `load`) в `TaskGroup`'ы
|
||||||
|
- Добавьте завершающую задачу `send_notification` на основе примера из раздела про алертинг (PythonOperator с `send_email_smtp` или другим механизмом уведомлений)
|
||||||
|
- Настройте для задачи уведомления `trigger_rule=TriggerRule.ALL_DONE`, чтобы уведомление отправлялось даже при частичных ошибках
|
||||||
|
- При желании вынесите функцию отправки письма в отдельный модуль `utils.py` и импортируйте её в DAG
|
||||||
|
|
||||||
|
**Цель задания:** Научиться группировать задачи с помощью TaskGroup и строить схему оповещений о статусе пайплайна.
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
## 🎯 Рекомендации по выполнению
|
## 🎯 Рекомендации по выполнению
|
||||||
|
|
||||||
1. **Начинайте с простых заданий** и постепенно переходите к сложным
|
1. **Начинайте с простых заданий** и постепенно переходите к сложным
|
||||||
@@ -257,6 +298,6 @@
|
|||||||
|
|
||||||
- 🟢 **Начальный уровень:** Выполнены задания для hello_world_dag и sql_basic_dag
|
- 🟢 **Начальный уровень:** Выполнены задания для hello_world_dag и sql_basic_dag
|
||||||
- 🟡 **Средний уровень:** Выполнены задания для file_operations_dag и data_processing_dag
|
- 🟡 **Средний уровень:** Выполнены задания для file_operations_dag и data_processing_dag
|
||||||
- 🟠 **Продвинутый уровень:** Выполнены все задания, включая branching_dag и error_handling_dag
|
- 🟠 **Продвинутый уровень:** Выполнены все задания, включая branching_dag, error_handling_dag и дополнительные задания по пулам, XCom, TaskGroup и алертингу
|
||||||
|
|
||||||
Удачи в изучении Apache Airflow! 🚀
|
Удачи в изучении Apache Airflow! 🚀
|
||||||
|
|||||||