diff --git a/01 - Введение в Airflow.md b/01 - Введение в Airflow.md index c5f3c54..d83cea8 100644 --- a/01 - Введение в Airflow.md +++ b/01 - Введение в Airflow.md @@ -1,9 +1,17 @@ # Почему Apache Airflow стал незаменимым инструментом для работы с данными -# Почему Apache Airflow стал незаменимым инструментом для работы с данными Обычно работа с автоматизацией процессов обработки информации начинается с ручного управления задачами. Например, в машинном обучении это может включать подготовку наборов данных, обучение моделей, анализ результатов и развертывание решений в рабочей среде. По мере роста команды и развития продукта эти процессы усложняются: увеличивается количество повторяющихся операций, появляются зависимости между задачами, и каждая из них приобретает всё большее значение для бизнеса. В результате формируется полноценный конвейер задач, требующий регулярного запуска. -![Введение в Apache Airflow](./_attachments/airflow_introduction.png) +```mermaid +flowchart LR + in[Данные] --> op1[Операция № 1] --> op2[Операция № 2] --> op3[Операция № 3] --> out[Данные] + + subgraph PIPE[Конвейер] + op1 + op2 + op3 + end +``` Аналогичная ситуация возникает при обработке данных: в определенный момент необходимо собрать актуальную информацию, преобразовать её и выполнить различные операции — создать витрину данных и сохранить в базу, обучить модель машинного обучения, подготовить отчет в Excel и разослать его по электронной почте. Вариантов множество, и для решения таких задач требуются специализированные инструменты. @@ -77,4 +85,4 @@ В этом материале вы познакомились с предпосылками появления инструментов управления процессами, подобных Airflow. Вы узнали, в каких сценариях Airflow наиболее эффективен, а в каких случаях стоит рассмотреть альтернативные решения, а также как этот инструмент применяется на практике в крупных российских и международных компаниях. -💡 Обратите внимание, что в данном модуле рассматривается Airflow версии 2.5. \ No newline at end of file +💡 Обратите внимание, что в данном модуле рассматривается Airflow версии 2.9. diff --git a/02 - Понятное введение в ключевые понятия Airflow.md b/02 - Понятное введение в ключевые понятия Airflow.md index ffadefa..a77bb43 100644 --- a/02 - Понятное введение в ключевые понятия Airflow.md +++ b/02 - Понятное введение в ключевые понятия Airflow.md @@ -12,7 +12,18 @@ Проще говоря, DAG — это упорядоченный набор задач, которые выполняются строго по расписанию и никогда не повторяются в рамках одного запуска. Создавая пайплайны для обработки данных, вы фактически создаете DAG. -![Пример направленного ациклического графа](_attachments/dag_example_directed_graph.png) +```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 зависит от успешного завершения всех предыдущих задач: -![DAG как последовательность задач](_attachments/dag_as_task_sequence.png) +```mermaid +flowchart LR + A["Task A"] --> B["Task B"] + A --> C["Task C"] + + B --> D["Task D"] + C --> D + + D --> E["Task E"] +``` Airflow позволяет создавать сложные сценарии: - Зависимости между разными DAG-ами с помощью TriggerDagRunOperator и ExternalTaskSensor - Условное выполнение задач в зависимости от результатов предыдущих шагов - Сложные ветвления и параллельные ветки выполнения -![Пример сложного DAG в Airflow](_attachments/complex_dag_example.png) +```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 +``` # Инструменты для выполнения: Операторы diff --git a/03 - Как устроен Airflow внутри.md b/03 - Как устроен Airflow внутри.md index 2643a44..0e108a4 100644 --- a/03 - Как устроен Airflow внутри.md +++ b/03 - Как устроен Airflow внутри.md @@ -8,7 +8,36 @@ Apache Airflow состоит из нескольких взаимосвязан На схеме ниже показано, как эти компоненты взаимодействуют между собой: -![](_attachments/airflow_architecture.png) +```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* Давайте подробно рассмотрим каждый компонент системы. diff --git a/04 - Знакомство с интерфейсом Airflow для начинающих.md b/04 - Знакомство с интерфейсом Airflow для начинающих.md index f5ae369..e82e54c 100644 --- a/04 - Знакомство с интерфейсом Airflow для начинающих.md +++ b/04 - Знакомство с интерфейсом Airflow для начинающих.md @@ -1,114 +1,88 @@ # Знакомство с интерфейсом 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): -![](_attachments/dag_list_status_indicators.png) +- **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-ами всё более-менее нормально или где-то горит?» -![](_attachments/dag_toggle_switches.png) +## Страница DAG: главное рабочее место -**Владелец процесса** — каждый пайплайн имеет ответственного владельца. Это особенно полезно в командной работе, когда несколько инженеров создают и поддерживают различные ETL-процессы. Владелец отвечает за мониторинг и корректную работу своего DAG. +Когда на главной странице (**DAGs View**) вы нажимаете на `DAG ID`, открывается страница конкретного пайплайна. -![](_attachments/dag_owner_field.png) +Верхняя часть страницы DAG: -**Статус выполнения** — цветные индикаторы с цифрами показывают количество и состояние последних запусков DAG: -- 🔴 Красный — завершено с ошибкой (failed) -- 🟡 Желтый — ожидает повторного запуска (retry) -- 🟢 Зеленый — выполняется в данный момент (running) -- 🟢 Тёмно-зеленый — успешно завершено (success) +* переключатель **Pause / Unpause**; +* кнопка **Trigger DAG** (ручной запуск); +* фильтр по дате, типу и состоянию запусков (**Run Type**, **Run State**, период по календарю); +* небольшой индикатор статусов задач (цветные ярлыки `running`, `failed`, `success` и т.п.). -![](_attachments/dag_status_colors.png) +> Для экспериментов удобно использовать стенд из папки `airflow-docker` в этом репозитории — там Airflow 2.9.2, и интерфейс будет выглядеть так же, как в учебнике. -**Расписание** — указывает, когда и с какой периодичностью запускается пайплайн. Используется формат cron, который может показаться сложным на первый взгляд. Для перевода cron-выражений в понятный формат рекомендуем использовать сервис [Crontab.guru](https://crontab.guru/). +Чуть ниже — горизонтальное меню вкладок: -![](_attachments/dag_schedule_field.png) +* **Details** + Краткое резюме DAG: количество задач, типы операторов, расписание, теги, статистика по запускам. Это удобная точка входа: «что это за DAG и как он в целом живёт». -**Последний запуск** — показывает дату и время самого свежего выполнения DAG, будь то автоматический запуск по расписанию или ручной запуск. +* **Graph** + Граф зависимостей задач. Здесь хорошо видно, какие задачи идут последовательно, какие — параллельно, где ветвления. + Клик по задаче открывает панель с действиями: **View Log**, **Clear**, **Mark Success / Mark Failed**, **Run** и др. -![](_attachments/dag_last_run_field.png) +* **Gantt** + Диаграмма Ганта для выбранного запуска DAG. Показывает, сколько времени заняла каждая задача и где они выполнялись параллельно. По ней удобно искать «бутылочные горлышки» — самые долгие шаги пайплайна. -**Статус задач** — детальная информация о последнем запуске: сколько задач находится в каждом статусе. Это помогает быстро оценить общее состояние пайплайна без необходимости погружаться в детали. +* **Run Duration** + История длительности запусков DAG. Помогает увидеть, не стали ли запуски в целом работать заметно дольше, и отследить, после какого изменения время выполнения выросло. -![](_attachments/dag_task_status_field.png) +* **Calendar** + Календарный вид истории запусков: по дням и месяцам видно, когда DAG запускался и как часто были ошибки. -**Быстрые действия** — в последнем столбце расположены кнопки для немедленного выполнения операций: запуск, обновление и удаление DAG. На практике этими кнопками пользуются редко. +* **Code** + Исходный код DAG, который сейчас задеплоен в Airflow. Быстрый способ проверить, что в среде действительно лежит та версия DAG, которую вы ждёте (и что изменения из Git уже подхватились). -Главная страница дает вам общее представление о состоянии всех ваших процессов. Но для детальной работы с конкретным пайплайном нужно перейти внутрь — просто кликните по названию интересующего DAG. +* **Audit Log** + Журнал действий по DAG: кто запускал, очищал задачи, менял состояние и т.д. Полезен, когда нужно понять, «кто и что нажал» перед тем, как всё сломалось. -![](_attachments/dag_click_to_open.png) +### Как работать с задачами (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. -Вы можете увидеть: -- Состав DAG и последовательность выполнения задач -- Тип оператора для каждой задачи -- Историю запусков в виде цветных квадратов напротив каждой задачи +Подробное описание всех экранов Airflow (с актуальными скриншотами) есть в официальной документации: [UI / Screenshots (Apache Airflow 2.9.3)][1]. -![](_attachments/tree_view_example.png) -![](_attachments/tree_view_zoomed.png) +Там же описаны **DAGs View**, **Grid View**, **Graph View**, **Gantt Chart**, **Task Duration**, **Landing Times**, **Code View**, **Audit Log** и другие разделы UI. -## Графическое представление (Graph View) +--- -Когда DAG содержит много задач, древовидное представление может быть неудобным. В таких случаях используйте вкладку Graph View — она показывает пайплайн в виде наглядного графа с четкими связями между задачами. - -![](_attachments/graph_view_example.png) -![](_attachments/graph_view_detailed.png) - -При клике на любую задачу открывается подробное окно с двумя основными разделами: - -### Информация о задаче -- **Просмотр логов** — переход к странице с полным выводом выполнения задачи (одна из самых часто используемых функций) -- **Детали выполнения** — подробная информация о конкретном запуске задачи - -### Управление задачей -Доступны четыре основных действия: -- **Запустить** — выполнить задачу немедленно -- **Очистить состояние** — сбросить статус задачи для повторного выполнения -- **Отметить как неудачную** — вручную установить статус ошибки -- **Отметить как успешную** — вручную установить статус успеха - -Каждое действие можно комбинировать с дополнительными опциями: -- **Игнорировать зависимости** — запуск без проверки зависимостей от других задач -- **Работать с прошлыми/будущими запусками** — применить действие ко всем запускам в определенном временном диапазоне -- **Влиять на связанные задачи** — применить действие к предыдущим (upstream) или последующим (downstream) задачам - -Наиболее популярная комбинация — **Downstream + Recursive + Clear**, которая сбрасывает текущую задачу и все зависящие от нее задачи в рамках одного запуска. - -## Анализ времени выполнения (Task Duration) - -Вкладка Task Duration автоматически строит графики на основе истории выполнения, показывая, сколько времени занимает каждая задача при каждом запуске. Это помогает выявлять узкие места и отслеживать изменения производительности. - -![](_attachments/task_duration_chart.png) - -## Диаграмма Ганта (Gantt) - -Диаграмма Ганта визуализирует распределение времени выполнения задач в рамках одного запуска DAG. Это отличный инструмент для определения самых ресурсоемких операций и планирования оптимизации. - -![](_attachments/gantt_chart_example.png) - -## Исходный код (Code) - -Вкладка Code отображает актуальный код DAG, который Airflow использует для выполнения. Это особенно полезно для проверки, что изменения из вашего Git-репозитория успешно загружены в систему и готовы к выполнению. - -![](_attachments/dag_code_view.png) - -Теперь вы знакомы с основными возможностями веб-интерфейса Airflow для мониторинга и управления вашими процессами обработки данных. Эти знания помогут вам эффективно работать с пайплайнами и быстро решать возникающие проблемы. \ No newline at end of file +[1]: https://airflow.apache.org/docs/apache-airflow/2.9.3/ui.html "UI / Screenshots — Airflow Documentation" diff --git a/05 - Основы построения DAG-файлов в Airflow.md b/05 - Основы построения DAG-файлов в Airflow.md index 4b47bfa..237532f 100644 --- a/05 - Основы построения DAG-файлов в Airflow.md +++ b/05 - Основы построения DAG-файлов в Airflow.md @@ -4,20 +4,63 @@ # Как устроен код DAG-файла -Как вы уже знаете из предыдущих уроков, DAG представляет собой граф вычислений, состоящий из отдельных задач (tasks). На уровне кода DAG — это обычный Python-файл, который описывает все задачи в рамках пайплайна и определяет последовательность их выполнения. +DAG представляет собой граф вычислений, состоящий из отдельных задач (tasks). На уровне кода DAG‑файл — это обычный Python‑скрипт, который Airflow регулярно импортирует, чтобы «увидеть» ваши пайплайны. -Любой DAG-файл состоит из нескольких ключевых компонентов: -- Импорт необходимых модулей и библиотек -- Настройка параметров и инициализация объекта DAG -- Создание отдельных задач с помощью операторов -- Определение порядка выполнения задач +Удобно мысленно разбивать любой DAG‑файл на четыре блока: +1. Импорт необходимых модулей и операторов +2. Настройка параметров и инициализация объекта `DAG` +3. Создание отдельных задач с помощью операторов +4. Определение порядка выполнения задач (задание зависимостей) -Порядок этих компонентов имеет значение, поскольку Airflow — это Python-библиотека, и обращение к еще не инициализированным объектам приведет к ошибкам выполнения. +Порядок этих блоков важен: как и в любом Python‑коде, нельзя обращаться к объектам, которые ещё не созданы. -Давайте рассмотрим простой пример DAG-файла и разберем его по частям. +Ниже приведён простой пример DAG‑файла. В следующих подразделах мы разберём каждый блок по отдельности, опираясь на этот пример. -![](_attachments/dag_file_structure.png) -*Пример структуры DAG-файла* +```python +# Секция импортов +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 -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 -t1 >> [t2, t3] # t2 и t3 запускаются одновременно после t1 -[t2, t3] >> t4 # t4 запускается после завершения t2 и t3 +t1 >> [t2, t3] # t2 и t3 стартуют после t1 +[t2, t3] >> t4 # t4 стартует после завершения и t2, и t3 ``` -Для сложных пайплайнов можно создавать цепочки любой сложности: +Можно комбинировать: + ```python 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 с более подробной настройкой параметров. Обратите внимание на дополнительные параметры, такие как количество повторных попыток ('retries'), задержка между попытками ('retry_delay'), а также настройки уведомлений по электронной почте. @@ -178,10 +243,9 @@ t1 >> t2 2. Создание агрегированной таблицы по регионам и категориям 3. Сохранение результата в базу данных PostgreSQL -Вот полный код нашего DAG: +Вот полный код нашего DAG, адаптированный под учебный стенд из папки `airflow-docker`: ```python -import os import datetime as dt import pandas as pd from airflow.models import DAG @@ -191,6 +255,10 @@ from airflow.operators.bash import BashOperator from sqlalchemy import create_engine +# Подключение к учебной базе PostgreSQL (postgres-training) +DB_URL = "postgresql://student:student@postgres-training:5432/training" + + # Базовые параметры DAG args = { 'owner': 'airflow', @@ -199,25 +267,30 @@ args = { 'retry_delay': dt.timedelta(minutes=1), } + def download_titanic_dataset(): + """Загрузка датасета Titanic и сохранение в базу""" url = 'https://web.stanford.edu/class/archive/cs/cs109/cs109.1166/stuff/titanic.csv' df = pd.read_csv(url) - engine = create_engine('postgresql+psycopg2://jovyan:jovyan@localhost:5432/de') + + engine = create_engine(DB_URL) df.to_sql('titanic', engine, index=False, if_exists='replace', schema='public') def pivot_dataset(): - engine = create_engine('postgresql+psycopg2://jovyan:jovyan@localhost:5432/de') + """Построение сводной таблицы и сохранение результата""" + engine = create_engine(DB_URL) titanic_df = pd.read_sql('select * from public.titanic', con=engine) df = titanic_df.pivot_table( - index=['Sex'], - columns=['Pclass'], - values='Name', - aggfunc='count' - ).reset_index() + index=['Sex'], + columns=['Pclass'], + values='Name', + aggfunc='count' + ).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_id='titanic_pivot', @@ -228,7 +301,7 @@ dag = DAG( # Начальная задача для логирования start = BashOperator( task_id='start', - bash_command='echo "Начинаем выполнение пайплайна! "', + bash_command='echo "Начинаем выполнение пайплайна!"', dag=dag, ) @@ -250,37 +323,23 @@ pivot_titanic_dataset = PythonOperator( 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 в учебной среде выполните следующие шаги: -1. Сохраните код в файл с расширением `.py` (например, `customer_analysis_dag.py`) - -Для тестирования нашего DAG в учебной среде выполните следующие шаги: - -1. Сохраните код в файл с расширением `.py` (например, `customer_analysis_dag.py`) - -2. Найдите запущенный контейнер с учебной средой Airflow: +1. Сохраните код в файл с расширением `.py` в папке стенда, например `airflow-docker/dags/titanic_pivot_dag.py`. +2. Убедитесь, что стенд запущен: ```bash -docker ps +cd airflow-docker +docker-compose up -d ``` +3. Подождите 30–60 секунд — Airflow автоматически обнаружит новый файл в папке `dags`. +4. Откройте веб-интерфейс Airflow (http://localhost:8080), найдите DAG `titanic_pivot` по идентификатору и запустите его. -3. Скопируйте файл в контейнер: -```bash -docker cp /путь/к/файлу/customer_analysis_dag.py [ID_КОНТЕЙНЕРА]:/lessons/dags/customer_analysis_dag.py -``` +После запуска вы можете посмотреть статус задач и логи в интерфейсе Airflow +(подробнее про это — в разделе про пользовательский интерфейс). -Например: -```bash -docker cp ~/Desktop/customer_analysis_dag.py 4dc4fbfedcf4:/lessons/dags/customer_analysis_dag.py -``` - -4. Подождите 30-60 секунд — Airflow автоматически обнаружит новый файл - -5. Найдите ваш DAG в интерфейсе Airflow через строку поиска и запустите его - -![Пример запуска DAG в интерфейсе Airflow](_attachments/dag_interface_example.png) - -В этом уроке вы изучили основную структуру DAG-файлов, создали свой первый рабочий пайплайн, научились загружать его в учебную среду и запускать для получения результатов. \ No newline at end of file +В этом уроке вы изучили основную структуру DAG-файлов, создали свой первый рабочий пайплайн, +научились загружать его в учебную среду и запускать для получения результатов. diff --git a/06 - Статусы задач в Airflow.md b/06 - Статусы задач в Airflow.md index 193aa28..79453ea 100644 --- a/06 - Статусы задач в Airflow.md +++ b/06 - Статусы задач в Airflow.md @@ -1,6 +1,4 @@ -# Статусы задач в Airflow - -# Статусы задач в Airflow — что означают цвета и как ими пользоваться +# Статусы задач в Airflow: как понимать и использовать ## Почему статусы задач так важны? @@ -8,56 +6,115 @@ ## Основные статусы, которые вы увидите каждый день -В интерфейсе Airflow каждая задача отображается определенным цветом. Вот что означают самые важные статусы: +В интерфейсе Airflow каждая задача подсвечивается цветом — по нему можно быстро понять, что с ней происходит. Для первых шагов достаточно запомнить несколько базовых статусов, которые удобно разделить на три группы. -![Статусы задач в интерфейсе Airflow](_attachments/task_status_interface.png) - -### Простое объяснение всех статусов - -Давайте разберем каждый статус простым языком: +**1. Всё хорошо** **🟢 Успешно (success)** — ваша задача выполнилась без ошибок. Это то, к чему мы стремимся! +**🟣 Пропущена (skipped)** — задача была намеренно пропущена (часто в ветвящихся пайплайнах). + +**2. Есть проблема** + **🔴 Ошибка (failed)** — что-то пошло не так. Задача упала, и вам нужно разбираться в коде. -**🟡 В очереди (queued)** — задача ждет своей очереди на выполнение. Это нормально, особенно если у вас много задач или мало ресурсов. +**3. Идёт работа или ожидание** **🔵 Выполняется (running)** — задача сейчас активно работает. Просто подождите немного. -**⚪ Нет статуса (none/no status)** — задача еще не готова к запуску, потому что не выполнены её зависимости. +**🟡 В очереди (queued)** — задача ждет своей очереди на выполнение. Это нормально, особенно если у вас много задач или мало ресурсов. **🟠 Запланирована (scheduled)** — все готово к запуску, Airflow вот-вот начнет выполнение. -**🟣 Пропущена (skipped)** — задача была намеренно пропущена (часто в ветвящихся пайплайнах). +**⚪ Нет статуса (none/no status)** — задача еще не готова к запуску, потому что не выполнены её зависимости. ### Специальные статусы (встречаются реже) -- **Ошибка в зависимости (upstream_failed)** — предыдущая задача упала, поэтому текущая даже не запускалась -- **Готова к повтору (up_for_retry)** — задача упала, но Airflow попробует запустить её снова (если настроены повторные попытки) -- **Завершена (shutdown)** — задачу принудительно остановили во время выполнения -- **Отложена (deferred)** — задача приостановлена и ждет внешнего события -- **Наблюдение (sensing)** — специальный статус для сенсоров, которые ждут определенных условий +- **Ошибка в зависимости (upstream_failed)** — предыдущая задача упала, поэтому текущая даже не запускалась. +- **Готова к повтору (up_for_retry)** — задача упала, но Airflow попробует запустить её снова (если настроены повторные попытки). +- **Завершена (shutdown)** — задачу принудительно остановили во время выполнения. +- **Отложена (deferred)** — задача приостановлена и ждет внешнего события. +- **Наблюдение (sensing)** — специальный статус для сенсоров, которые ждут определенных условий. ## Как задача проходит свой путь: пошагово Представьте, что у вас есть простая задача. Вот как она проходит свой жизненный цикл: 1. **Создание** → Статус: "Нет статуса" - Airflow создает задачу, но еще не может её запустить + Airflow создает задачу, но еще не может её запустить. 2. **Готовность** → Статус: "Запланирована" - Все зависимости выполнены, задача готова к работе + Все зависимости выполнены, задача готова к работе. 3. **Ожидание** → Статус: "В очереди" - Задача ждет свободного рабочего места + Задача ждет свободного рабочего места. 4. **Работа** → Статус: "Выполняется" - Задача активно выполняется + Задача активно выполняется. 5. **Завершение** → Статус: "Успешно" Всё прошло отлично! -![Жизненный цикл задачи в Airflow](_attachments/task_lifecycle_detailed.png) +На диаграмме ниже показаны те же этапы, но уже с привязкой к внутренним компонентам 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 @@ - **Остальные статусы** вы будете изучать по мере необходимости - **Цвета в интерфейсе** — это быстрый способ понять состояние вашего пайплайна -Помните: понимание статусов задач — это как научиться читать дорожные знаки. Сначала кажется много информации, но со временем это становится второй натурой! \ No newline at end of file +Помните: понимание статусов задач — это как научиться читать дорожные знаки. Сначала кажется много информации, но со временем это становится второй натурой! + +> Полный список возможных состояний задач (TaskInstanceState) и их классификацию на терминальные и промежуточные можно посмотреть в официальной документации Airflow: +> https://airflow.apache.org/docs/apache-airflow/2.9.3/_api/airflow/utils/state/index.html + diff --git a/07 - Гибкие шаблоны и настройки в Airflow.md b/07 - Шаблоны, переменные и подключения в Airflow.md similarity index 60% rename from 07 - Гибкие шаблоны и настройки в Airflow.md rename to 07 - Шаблоны, переменные и подключения в Airflow.md index 12139b6..618fc87 100644 --- a/07 - Гибкие шаблоны и настройки в Airflow.md +++ b/07 - Шаблоны, переменные и подключения в Airflow.md @@ -1,4 +1,5 @@ -# Гибкие шаблоны и настройки в Airflow +# Шаблоны, переменные и подключения в Airflow + В этом материале вы познакомитесь с мощными инструментами Airflow для создания гибких и переиспользуемых пайплайнов: динамическими шаблонами, безопасными переменными и централизованными подключениями к внешним системам. # Динамические шаблоны Airflow (на основе Jinja) @@ -18,30 +19,30 @@ ... dag = DAG( - dag_id="dynamic_templates_example", - schedule_interval="*/10 * * * *", + dag_id="template_example", + schedule_interval="*/15 * * * *", 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 ``` -**Использование шаблонов в идентификаторах задач для лучшей отслеживаемости:** +**Использование шаблонов для лучшей отслеживаемости задач:** -В этом примере создается DAG с идентификатором "template_tracking_example", который запускается каждые 20 минут. Вторая задача использует шаблон {{ ds }} в своем идентификаторе, что позволяет легко идентифицировать задачу по дате выполнения. +В этом примере создается DAG с идентификатором "template_tracking_example", который запускается каждые 20 минут. Вторая задача использует шаблон {{ ds }} в команде bash, чтобы явно указывать дату обработки в логах. ```python ... dag = DAG( - dag_id="dynamic_templates_example", - schedule_interval="*/10 * * * *", + dag_id="template_tracking_example", + schedule_interval="*/20 * * * *", default_args=default_args ) t1 = BashOperator(task_id="show_date", bash_command="echo {{ ds }}") -t2 = BashOperator(task_id=f"process_for_{{ ds }}", bash_command="echo {{ ds }}") +t2 = BashOperator(task_id="process_for_date", bash_command="echo Processing for {{ ds }}") t1 >> t2 ``` @@ -50,22 +51,52 @@ t1 >> t2 Переменные Airflow представляют собой пары "ключ-значение", хранящиеся в метадатабазе системы. Они идеально подходят для хранения конфигурационных параметров, таких как пути к скриптам, имена таблиц или другие настройки, которые должны быть доступны в разных DAG. -Управление переменными осуществляется через веб-интерфейс Airflow (Admin → Variables), где можно: -- Создавать и редактировать пары ключ-значение вручную -- Импортировать настройки из JSON-файлов -- Использовать командную строку Airflow +Управление переменными осуществляется через веб-интерфейс Airflow (раздел **Admin → Variables**). Через этот раздел можно: -![Управление переменными через UI](_attachments/variables_management_ui.gif) +- создавать и редактировать пары «ключ-значение» вручную; +- импортировать набор переменных из 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 с идентификатором "variable_example", который использует переменную 'data_storage_path', предварительно сохраненную в Airflow. Значение переменной извлекается с помощью Variable.get() и используется в команде bash для указания пути к данным. + ```python from airflow import DAG from airflow.operators.bash import BashOperator -from airflow.operators.dummy import DummyOperator from airflow.models import Variable from datetime import datetime @@ -108,7 +139,7 @@ task = BashOperator( В приведенном примере используется PostgresOperator для создания таблицы в базе данных. Вместо использования подключения по умолчанию, явно указывается подключение с идентификатором 'my_postgres_conn'. ```python -from airflow.operators.postgres_operator import PostgresOperator +from airflow.providers.postgres.operators.postgres import PostgresOperator create_table = PostgresOperator( 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). -![Настройка подключения к PostgreSQL](_attachments/postgres_connection_setup.gif) +Управление подключениями доступно через интерфейс 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. -Более подробную информацию о настройке подключений можно найти в [документации 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 @@ -139,4 +212,4 @@ create_table = PostgresOperator( - Требуется ли маскировка конфиденциальных значений? - Используются ли в логике даты или временные метки, которые можно заменить на шаблоны? -Ответы на эти вопросы помогут вам создавать надежные, безопасные и легко поддерживаемые пайплайны в Airflow. \ No newline at end of file +Ответы на эти вопросы помогут вам создавать надежные, безопасные и легко поддерживаемые пайплайны в Airflow. diff --git a/08 - Управление временем.md b/08 - Управление временем.md index 0e687c0..a5722c5 100644 --- a/08 - Управление временем.md +++ b/08 - Управление временем.md @@ -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 установлена одна попытка повторного запуска при ошибках. @@ -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 с исторической датой начала, Airflow предлагает мощный механизм автоматического пересчета пропущенных периодов через параметр `catchup`. - ```python dag = DAG( dag_id="historical_data_processing", @@ -61,7 +59,7 @@ dag = DAG( При `catchup=True` система автоматически выполнит все пропущенные запуски от указанной даты начала до текущего момента. Это особенно полезно при первом запуске DAG для обработки накопившихся исторических данных. -Если ваш бизнес-сценарий не требует пересчета истории или вы хотите начать обработку только с текущего периода, установите `catchup=False`. В этом случае Airflow будет планировать только ближайшие запуски согласно расписанию. +По умолчанию для DAG с расписанием параметр `catchup` включен (`True`), поэтому при первом деплое с исторической `start_date` можно неожиданно получить большое количество запусков. Если ваш бизнес-сценарий не требует пересчета истории или вы хотите начать обработку только с текущего периода, установите `catchup=False`. В этом случае Airflow будет планировать только ближайшие запуски согласно расписанию. ## Ручная перезаливка данных (backfill) @@ -82,7 +80,7 @@ the_main_dag Вы можете: - Настроить глобальную временную зону в конфигурационном файле Airflow -- Указать временную зону явно при инициализации DAG через параметр `tz` +- Указать временную зону для DAG через параметр `timezone` или передать в `start_date` объект с явной таймзоной (например, созданный через `pendulum.datetime(..., tz="UTC")`) Однако рекомендуется придерживаться UTC во всех расчетах и преобразовывать временные метки только при выводе результатов для конечных пользователей. Это минимизирует ошибки и упрощает отладку. @@ -93,4 +91,4 @@ the_main_dag 3. **Планируйте запуски с учетом UTC**, особенно если ваша команда работает в разных часовых поясах 4. **Документируйте временные зависимости** в коде DAG для других разработчиков -Понимание временных механизмов Airflow — ключ к созданию надежных и предсказуемых пайплайнов. Правильная настройка расписаний и управление историческими данными позволяют автоматизировать сложные бизнес-процессы без ручного вмешательства. \ No newline at end of file +Понимание временных механизмов Airflow — ключ к созданию надежных и предсказуемых пайплайнов. Правильная настройка расписаний и управление историческими данными позволяют автоматизировать сложные бизнес-процессы без ручного вмешательства. diff --git a/09 - Продвинутые возможности Airflow.md b/09 - Пулы, XCom, TaskGroup и алертинг в Airflow.md similarity index 79% rename from 09 - Продвинутые возможности Airflow.md rename to 09 - Пулы, XCom, TaskGroup и алертинг в Airflow.md index f2df916..8e1443b 100644 --- a/09 - Продвинутые возможности Airflow.md +++ b/09 - Пулы, XCom, TaskGroup и алертинг в Airflow.md @@ -1,6 +1,6 @@ -# Продвинутые возможности Airflow для начинающих специалистов +# Пулы, XCom, TaskGroup и алертинг в Airflow -В этом материале мы познакомимся с расширенными функциями Airflow, которые помогут вам решать нетривиальные задачи при работе со сложными пайплайнами. +В этом материале мы познакомимся с управлением ресурсами через пулы задач, обменом данными между задачами через XCom, группировкой задач с помощью TaskGroup и настройкой системы оповещений (алертинга) в Airflow. Начнем с улучшения нашего базового пайплайна — проведем рефакторинг кода для лучшей читаемости и поддержки. @@ -12,11 +12,18 @@ Когда планировщик обнаруживает, что наступило время выполнения DAG, он запускает задачу согласно заданной последовательности. При этом задача занимает один слот в пуле и освобождает его после завершения. -Создать новый пул задач и установить его размер можно через веб-интерфейс Airflow: +Создать новый пул задач и установить его размер можно через веб-интерфейс Airflow. Рассмотрим пример пула `data_processing_pool` для тяжёлых задач: -![image](_attachments/advanced_feature_demo_1.gif) +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 и другие. Контекст можно найти следующим образом: -![image](_attachments/advanced_feature_demo_3.gif){: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 }}`. -![image](_attachments/advanced_feature_demo_2.gif) +В 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 предназначен для обмена небольшими сообщениями. Данные проходят сериализацию/десериализацию при чтении и записи в таблицу. @@ -73,7 +101,7 @@ BashOperator( Для передачи больших объемов данных используйте внешние средства: файловую систему, базы данных (чаще всего PostgreSQL) или очереди сообщений (например, Kafka). -Механизм XCom похож на работу функций в Python. Многие операторы по умолчанию возвращают результат выполнения задачи. За это отвечает параметр `do_xcom_push`, который по умолчанию равен `True`. +Механизм XCom похож на работу функций в Python. Многие операторы (например, `PythonOperator`) по умолчанию возвращают результат выполнения задачи. За это отвечает параметр `do_xcom_push`, который во многих случаях равен `True` по умолчанию. Чтобы прочитать сообщения из XCom, используйте метод `xcom_pull` в контексте задачи: @@ -149,17 +177,44 @@ with TaskGroup("data_loading") as loading_group: Обратите внимание: при использовании TaskGroup последовательность задач указывается внутри группы после объявления всех задач, а в конце DAG описывается последовательность выполнения самих групп. -Визуально в интерфейсе Airflow это выглядит так: +Ниже приведена упрощённая схема зависимостей между тремя группами задач: -![image](_attachments/advanced_feature_diagram_17.png) +```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 -![image](_attachments/advanced_feature_animation_23.gif) + %% 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 — это удобный способ логической группировки задач, который помогает упростить код и представить сложные пайплайны более компактно. @@ -167,7 +222,7 @@ TaskGroup — это удобный способ логической групп Алертинг — один из ключевых компонентов системы оркестрации, так как важно своевременно получать уведомления об ошибках для их оперативного анализа и решения. -По умолчанию в Airflow настроена отправка уведомлений на электронную почту. При создании DAG указываются email-адреса, на которые будут отправляться сообщения. С помощью параметров можно настроить различные сценарии оповещений. +В Airflow есть встроенная поддержка отправки уведомлений на электронную почту (при условии, что в конфигурации настроен SMTP-сервер). При создании DAG указываются email-адреса, на которые будут отправляться сообщения. С помощью параметров можно настроить различные сценарии оповещений. Давайте модифицируем наш первый DAG так, чтобы получать уведомления на почту при возникновении ошибок. При этом настроим перезапуск задач в случае неудачи (например, 2 попытки), но без уведомлений о самих перезапусках: @@ -208,14 +263,16 @@ dag = DAG( from datetime import datetime, timedelta, timezone import dateutil from airflow.utils.email import send_email_smtp +from airflow.operators.python import get_current_context MAIL_LIST = [ "email_1@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. + context = get_current_context() calculation_start = dateutil.parser.isoparse(dagrun_begin_time.replace('Z', '+00:00')) calculation_end = datetime.now(timezone.utc) duration = str(timedelta(seconds=(calculation_end - calculation_start).seconds)) @@ -241,6 +298,7 @@ def notify_email(calculation_dt: str, dagrun_begin_time, **context): ```python from airflow.operators.python import PythonOperator +from airflow.utils.trigger_rule import TriggerRule from utils import notify_email 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( task_id="send_email_notification", python_callable=notify_email, - provide_context=True, dag=dag, trigger_rule=TriggerRule.ALL_DONE, op_args=[LOCAL_CALCULATION_DT, DAG_RUN_BEGIN_TIME], @@ -267,6 +324,7 @@ email_notification_task = PythonOperator( В этом примере используется переменная 'database_connection_string', предварительно созданная в интерфейсе Airflow, для хранения строки подключения к базе данных. Это позволяет избежать жесткого кодирования конфиденциальной информации в коде DAG и упрощает настройку подключения для разных окружений. ```python +import os import datetime as dt import pandas as pd 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) def load_customer_data(): - url = 'https://example.com/customer_data.csv' - df = pd.read_csv(url) + file_path = get_file_path('customer_data.csv') + df = pd.read_csv(file_path) engine = create_engine(DATABASE_URL) df.to_sql('customers', engine, index=False, if_exists='replace', schema='staging') @@ -351,6 +409,7 @@ start_task >> data_processing Итоговый DAG будет выглядеть следующим образом: ```python +import os import datetime as dt import pandas as pd 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) def load_customer_data(): - url = 'https://example.com/customer_data.csv' - df = pd.read_csv(url) + file_path = get_file_path('customer_data.csv') + df = pd.read_csv(file_path) engine = create_engine(DATABASE_URL) df.to_sql('customers', engine, index=False, if_exists='replace', schema='staging') @@ -426,4 +485,4 @@ start_task >> data_processing - Логическая группировка задач с TaskGroup - Настройка системы оповещений -Помните: не стоит использовать все доступные функции сразу. Выбирайте инструменты последовательно и находите оптимальный набор возможностей под конкретную задачу. \ No newline at end of file +Помните: не стоит использовать все доступные функции сразу. Выбирайте инструменты последовательно и находите оптимальный набор возможностей под конкретную задачу. diff --git a/AGENTS.md b/AGENTS.md index 51c5270..994938c 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -1,6 +1,6 @@ # 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 - Root `01-09 *.md`: step-by-step articles (RU). @@ -41,3 +41,6 @@ Prerequisite: Docker + Docker Compose. ## Security & Configuration Tips - 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. + +## Hints +- При работе под Windows для работы с командной строкой используй PowerShell diff --git a/README.md b/README.md index d3cfc60..66660a9 100644 --- a/README.md +++ b/README.md @@ -11,7 +11,7 @@ Apache Airflow — это мощный оркестратор рабочих п - Хотите освоить свой первый ETL-инструмент - Нуждаетесь в практическом руководстве с понятными объяснениями -Мы начнем с основ и постепенно перейдем к продвинутым возможностям, чтобы вы могли уверенно использовать Airflow в реальных проектах. Курс ориентирован на версию Airflow 2.5 и использует практический подход с множеством примеров и визуальных материалов. +Мы начнем с основ и постепенно перейдем к продвинутым возможностям, чтобы вы могли уверенно использовать Airflow в реальных проектах. Курс ориентирован на версию Airflow 2.9 и использует практический подход с множеством примеров и визуальных материалов. ## 🎯 Что вы узнаете @@ -21,7 +21,7 @@ Apache Airflow — это мощный оркестратор рабочих п - **Интерфейс**: как эффективно использовать веб-интерфейс для мониторинга и управления - **Практическое создание 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-шаблонов - Параметризация DAG -- Глобальные переменные и соединения +- Глобальные переменные и подключения ### [08. Управление временем в Airflow](08%20-%20Управление%20временем.md) - Расписания и интервалы запуска - Работа с временными зонами - Понимание execution_date и других временных концепций -### [09. Продвинутые возможности Airflow](09%20-%20Продвинутые%20возможности%20Airflow.md) -- Расширенные операторы и сенсоры -- Обработка ошибок и повторные попытки -- Оптимизация и масштабирование DAG +### [09. Пулы, XCom, TaskGroup и алертинг в Airflow](09%20-%20Пулы,%20XCom,%20TaskGroup%20и%20алертинг%20в%20Airflow.md) +- Управление ресурсами с помощью пулов задач +- Обмен данными между задачами через XCom +- Логическая группировка задач с TaskGroup +- Настройка системы оповещений и уведомлений ## 🛠️ Практический стенд @@ -102,4 +103,4 @@ Apache Airflow — это мощный оркестратор рабочих п --- -*Этот учебник создан для начинающих специалистов в области данных и инженерии. Все материалы ориентированы на практическое применение и пошаговое освоение Apache Airflow.* \ No newline at end of file +*Этот учебник создан для начинающих специалистов в области данных и инженерии. Все материалы ориентированы на практическое применение и пошаговое освоение Apache Airflow.* diff --git a/_attachments/advanced_feature_animation_23.gif b/_attachments/advanced_feature_animation_23.gif deleted file mode 100644 index 8ae977f..0000000 Binary files a/_attachments/advanced_feature_animation_23.gif and /dev/null differ diff --git a/_attachments/advanced_feature_demo_1.gif b/_attachments/advanced_feature_demo_1.gif deleted file mode 100644 index 62f991d..0000000 Binary files a/_attachments/advanced_feature_demo_1.gif and /dev/null differ diff --git a/_attachments/advanced_feature_demo_2.gif b/_attachments/advanced_feature_demo_2.gif deleted file mode 100644 index 2e205bf..0000000 Binary files a/_attachments/advanced_feature_demo_2.gif and /dev/null differ diff --git a/_attachments/advanced_feature_demo_3.gif b/_attachments/advanced_feature_demo_3.gif deleted file mode 100644 index 672762f..0000000 Binary files a/_attachments/advanced_feature_demo_3.gif and /dev/null differ diff --git a/_attachments/advanced_feature_diagram_17.png b/_attachments/advanced_feature_diagram_17.png deleted file mode 100644 index 6cee324..0000000 Binary files a/_attachments/advanced_feature_diagram_17.png and /dev/null differ diff --git a/_attachments/airflow_architecture.png b/_attachments/airflow_architecture.png deleted file mode 100644 index d799d4d..0000000 Binary files a/_attachments/airflow_architecture.png and /dev/null differ diff --git a/_attachments/airflow_introduction.png b/_attachments/airflow_introduction.png deleted file mode 100644 index aa74d7f..0000000 Binary files a/_attachments/airflow_introduction.png and /dev/null differ diff --git a/_attachments/complex_dag_example.png b/_attachments/complex_dag_example.png deleted file mode 100644 index 7d2cdf0..0000000 Binary files a/_attachments/complex_dag_example.png and /dev/null differ diff --git a/_attachments/dag_as_task_sequence.png b/_attachments/dag_as_task_sequence.png deleted file mode 100644 index 5d24927..0000000 Binary files a/_attachments/dag_as_task_sequence.png and /dev/null differ diff --git a/_attachments/dag_click_to_open.png b/_attachments/dag_click_to_open.png deleted file mode 100644 index 99675b7..0000000 Binary files a/_attachments/dag_click_to_open.png and /dev/null differ diff --git a/_attachments/dag_code_view.png b/_attachments/dag_code_view.png deleted file mode 100644 index 47d965b..0000000 Binary files a/_attachments/dag_code_view.png and /dev/null differ diff --git a/_attachments/dag_example_directed_graph.png b/_attachments/dag_example_directed_graph.png deleted file mode 100644 index 6826a5d..0000000 Binary files a/_attachments/dag_example_directed_graph.png and /dev/null differ diff --git a/_attachments/dag_file_structure.png b/_attachments/dag_file_structure.png deleted file mode 100644 index 90e5b22..0000000 Binary files a/_attachments/dag_file_structure.png and /dev/null differ diff --git a/_attachments/dag_interface_example.png b/_attachments/dag_interface_example.png deleted file mode 100644 index a65fa64..0000000 Binary files a/_attachments/dag_interface_example.png and /dev/null differ diff --git a/_attachments/dag_last_run_field.png b/_attachments/dag_last_run_field.png deleted file mode 100644 index 8124565..0000000 Binary files a/_attachments/dag_last_run_field.png and /dev/null differ diff --git a/_attachments/dag_list_status_indicators.png b/_attachments/dag_list_status_indicators.png deleted file mode 100644 index 702f2cc..0000000 Binary files a/_attachments/dag_list_status_indicators.png and /dev/null differ diff --git a/_attachments/dag_owner_field.png b/_attachments/dag_owner_field.png deleted file mode 100644 index ec7d3a6..0000000 Binary files a/_attachments/dag_owner_field.png and /dev/null differ diff --git a/_attachments/dag_schedule_field.png b/_attachments/dag_schedule_field.png deleted file mode 100644 index 2b1fe3f..0000000 Binary files a/_attachments/dag_schedule_field.png and /dev/null differ diff --git a/_attachments/dag_status_colors.png b/_attachments/dag_status_colors.png deleted file mode 100644 index 7e08c14..0000000 Binary files a/_attachments/dag_status_colors.png and /dev/null differ diff --git a/_attachments/dag_task_status_field.png b/_attachments/dag_task_status_field.png deleted file mode 100644 index de6f828..0000000 Binary files a/_attachments/dag_task_status_field.png and /dev/null differ diff --git a/_attachments/dag_toggle_switches.png b/_attachments/dag_toggle_switches.png deleted file mode 100644 index 40971bb..0000000 Binary files a/_attachments/dag_toggle_switches.png and /dev/null differ diff --git a/_attachments/gantt_chart_example.png b/_attachments/gantt_chart_example.png deleted file mode 100644 index ea713b4..0000000 Binary files a/_attachments/gantt_chart_example.png and /dev/null differ diff --git a/_attachments/graph_view_detailed.png b/_attachments/graph_view_detailed.png deleted file mode 100644 index b323b49..0000000 Binary files a/_attachments/graph_view_detailed.png and /dev/null differ diff --git a/_attachments/graph_view_example.png b/_attachments/graph_view_example.png deleted file mode 100644 index fc94612..0000000 Binary files a/_attachments/graph_view_example.png and /dev/null differ diff --git a/_attachments/postgres_connection_setup.gif b/_attachments/postgres_connection_setup.gif deleted file mode 100644 index f1c8219..0000000 Binary files a/_attachments/postgres_connection_setup.gif and /dev/null differ diff --git a/_attachments/task_duration_chart.png b/_attachments/task_duration_chart.png deleted file mode 100644 index 4f7888c..0000000 Binary files a/_attachments/task_duration_chart.png and /dev/null differ diff --git a/_attachments/task_lifecycle_detailed.png b/_attachments/task_lifecycle_detailed.png deleted file mode 100644 index f611d90..0000000 Binary files a/_attachments/task_lifecycle_detailed.png and /dev/null differ diff --git a/_attachments/task_status_interface.png b/_attachments/task_status_interface.png deleted file mode 100644 index 281aa83..0000000 Binary files a/_attachments/task_status_interface.png and /dev/null differ diff --git a/_attachments/tree_view_example.png b/_attachments/tree_view_example.png deleted file mode 100644 index 3a3bc03..0000000 Binary files a/_attachments/tree_view_example.png and /dev/null differ diff --git a/_attachments/tree_view_zoomed.png b/_attachments/tree_view_zoomed.png deleted file mode 100644 index a50655a..0000000 Binary files a/_attachments/tree_view_zoomed.png and /dev/null differ diff --git a/_attachments/variables_management_ui.gif b/_attachments/variables_management_ui.gif deleted file mode 100644 index d563b58..0000000 Binary files a/_attachments/variables_management_ui.gif and /dev/null differ diff --git a/airflow-docker/AGENTS.md b/airflow-docker/AGENTS.md index 6803d8e..ac29e1f 100644 --- a/airflow-docker/AGENTS.md +++ b/airflow-docker/AGENTS.md @@ -7,9 +7,8 @@ This file provides guidance for AI agents working on the educational Airflow set ``` airflow-docker/ -├── docker-compose.yml # Docker configuration -├── .env # Hardcoded environment variables -├── dags/ # DAG files for learning +├── docker-compose.yml # Docker configuration and environment +├── dags/ # DAG files for learning │ ├── hello_world_dag.py # Basic Python operators │ ├── sql_basic_dag.py # SQL operations │ ├── file_operations_dag.py # File processing @@ -118,8 +117,7 @@ docker-compose logs postgres-metadata ## 📚 Key Files to Examine ### Configuration Files -- [`docker-compose.yml`](docker-compose.yml) - Main Docker configuration -- [`.env`](.env) - Environment variables +- [`docker-compose.yml`](docker-compose.yml) - Main Docker configuration and environment variables - [`requirements.txt`](requirements.txt) - Python dependencies ### Sample Data Files @@ -184,4 +182,4 @@ docker-compose logs postgres-metadata ### For New DAGs - Include comprehensive docstring in Russian - Describe learning objectives clearly -- Provide step-by-step task explanations \ No newline at end of file +- Provide step-by-step task explanations diff --git a/airflow-docker/dag-specifications.md b/airflow-docker/dag-specifications.md index abb3342..5854063 100644 --- a/airflow-docker/dag-specifications.md +++ b/airflow-docker/dag-specifications.md @@ -131,6 +131,27 @@ Process data based on file type or data quality - `success_handler`: On success 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 ### customers.csv @@ -214,6 +235,12 @@ CREATE TABLE enrollments ( - ✅ Use parameters and templates - ✅ 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 ### Problem: DAG not appearing in UI diff --git a/airflow-docker/data/.gitkeep b/airflow-docker/data/.gitkeep old mode 100644 new mode 100755 diff --git a/airflow-docker/data/input/.gitkeep b/airflow-docker/data/input/.gitkeep old mode 100644 new mode 100755 diff --git a/airflow-docker/data/output/.gitkeep b/airflow-docker/data/output/.gitkeep old mode 100644 new mode 100755 diff --git a/airflow-docker/educational-tasks.md b/airflow-docker/educational-tasks.md index 68d84c2..b1a33e9 100644 --- a/airflow-docker/educational-tasks.md +++ b/airflow-docker/educational-tasks.md @@ -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. **Начинайте с простых заданий** и постепенно переходите к сложным @@ -257,6 +298,6 @@ - 🟢 **Начальный уровень:** Выполнены задания для hello_world_dag и sql_basic_dag - 🟡 **Средний уровень:** Выполнены задания для file_operations_dag и data_processing_dag -- 🟠 **Продвинутый уровень:** Выполнены все задания, включая branching_dag и error_handling_dag +- 🟠 **Продвинутый уровень:** Выполнены все задания, включая branching_dag, error_handling_dag и дополнительные задания по пулам, XCom, TaskGroup и алертингу -Удачи в изучении Apache Airflow! 🚀 \ No newline at end of file +Удачи в изучении Apache Airflow! 🚀