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/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! 🚀