Новые задания по pool, xcom, taskgroup
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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! 🚀
|
||||
Reference in New Issue
Block a user