Files
ddadmin ab1b61f396 docs(airflow): обновлены практические задания
- Зачем:
  - упражнения должны быть понятнее и соответствовать ожидаемому уровню сложности.
- Что:
  - добавлены подсказки и описания подводных камней.
  - упрощены задания по отчетам, мониторингу и логированию.
- Проверка:
  - git diff --check.
2026-07-26 23:58:10 +03:00

479 lines
36 KiB
Markdown
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# Учебные задания по Apache Airflow
Этот документ содержит практические задания для закрепления знаний по Apache Airflow. Каждое задание соответствует одному из учебных DAG'ов и направлено на лучшее понимание конкретных концепций.
## Общие инструкции
Перед выполнением заданий:
1. Убедитесь, что стенд Airflow запущен: `docker-compose up -d`
2. Проверьте доступность интерфейса: http://localhost:8080
3. Ознакомьтесь с соответствующим DAG'ом в интерфейсе Airflow
4. Откройте исходный код DAG'а и прочитайте описание в начале файла
---
## Задания для hello_world_dag.py
**Цель:** Освоить базовые концепции Airflow — создание задач, зависимости, операторы.
### Задание 1: Модификация существующих задач
**Сложность:** Начальная
**Задача:**
- Откройте файл `hello_world_dag.py`
- Измените текст в функциях `print_hello()`, `print_date()`, `print_goodbye()` на русский язык
- Добавьте новую задачу, которая выводит текущую погоду (можно использовать фиктивные данные)
**Подсказка:**
- Новая задача создаётся по аналогии с `date_task` — нужна Python-функция и `PythonOperator` для неё
- Не забудьте встроить задачу в цепочку зависимостей (`>>`)
- После сохранения файла DAG подхватится автоматически — запустите его вручную через UI и проверьте логи
**Цель задания:** Понять структуру DAG, научиться добавлять и модифицировать задачи.
### Задание 2: Изменение расписания
**Сложность:** Начальная
**Задача:**
- Измените `schedule_interval` с ежедневного на еженедельное выполнение
- Добавьте параметр `max_active_runs` для ограничения одновременных запусков
- Проверьте изменения в интерфейсе Airflow
**Подсказка:**
- Расписание задаётся в параметрах `DAG(...)` — ищите `schedule_interval`
- Можно использовать `timedelta(weeks=1)` или cron-выражение (например, `'0 0 * * 1'`)
- `max_active_runs` — тоже параметр `DAG()`
**Цель задания:** Освоить настройку расписания и параметров выполнения DAG.
### Задание 3: Добавление BashOperator
**Сложность:** Средняя
**Задача:**
- Добавьте задачу с `BashOperator`, которая создает текстовый файл в папке `/opt/airflow/data/output/`
- Настройте зависимости так, чтобы эта задача выполнялась между `date_task` и `end_task`
- Убедитесь, что файл создается при каждом запуске DAG
**Подсказка:**
- Импорт `BashOperator` уже есть в файле — посмотрите шапку
- `BashOperator` принимает параметр `bash_command` со строкой shell-команды
- Для записи в файл подойдёт обычный `echo "текст" > /path/to/file.txt`
- Проверить результат можно: `docker-compose exec airflow-webserver cat /opt/airflow/data/output/<ваш_файл>`
**Цель задания:** Научиться работать с разными типами операторов.
---
## Задания для sql_basic_dag.py
**Цель:** Освоить работу с базами данных через Airflow.
### Задание 1: Расширение структуры таблицы
**Сложность:** Средняя
**Задача:**
- Добавьте в таблицу `students_sample` новые поля: `email`, `phone`, `registration_date`
- Модифицируйте SQL-запросы для работы с новой структурой
- Добавьте задачу для обновления существующих записей
**Подсказка:**
- Новые колонки добавляются в переменную `create_table_sql` — подберите подходящие типы (VARCHAR, DATE и т.д.)
- Не забудьте обновить `insert_data_sql` — количество значений в INSERT должно совпадать с количеством колонок
- Для задачи UPDATE создайте новую SQL-переменную и новый `PostgresOperator` по аналогии с существующими
- **Подводный камень:** `TRUNCATE TABLE` не изменяет структуру. Чтобы новые колонки появились, замените его на `DROP TABLE IF EXISTS`
**Цель задания:** Научиться работать с миграциями схемы базы данных.
### Задание 2: Создание отчетов
**Сложность:** Средняя
**Задача:**
- Добавьте задачу, которая создает сводный отчет по данным студентов
- Отчет должен содержать: количество студентов, средний возраст, распределение по возрасту
- Сохраните отчет в файл в папке `/opt/airflow/data/output/`
**Подсказка:**
- Здесь понадобится `PythonOperator`, а не `PostgresOperator` — потому что нужно получить результат SQL и записать в файл
- Для подключения к БД из Python-функции используйте `PostgresHook` (импорт: `from airflow.providers.postgres.hooks.postgres import PostgresHook`)
- У хука есть метод `get_conn()`, который возвращает обычное psycopg2-соединение — дальше работайте через `cursor.execute()` и `fetchone()`
- Задачу поставьте перед `drop_table` в цепочке зависимостей
**Цель задания:** Освоить создание аналитических отчетов в DAG'ах.
### Задание 3: Работа с соединениями
**Сложность:** Средняя
**Задача:**
- Создайте новое соединение к базе данных через Airflow UI
- Модифицируйте DAG для использования нового соединения
- Добавьте обработку ошибок подключения к БД
**Подсказка:**
- Соединения создаются в UI: Admin → Connections → "+"
- Параметры для нашей учебной БД: Host = `postgres-training`, Schema = `training`, Login = `student`, Password = `student`, Port = `5432`
- В DAG-файле соединение указывается через `postgres_conn_id` — замените значение на ваш новый Conn Id
- Для проверки ошибок попробуйте указать несуществующее соединение и посмотрите, что покажут логи
**Цель задания:** Научиться управлять соединениями с внешними системами.
---
## Задания для file_operations_dag.py
**Цель:** Освоить работу с файлами и данными в Airflow.
### Задание 1: Модификация генерации данных
**Сложность:** Средняя
**Задача:**
- Добавьте новые поля в генерируемые данные: `department`, `experience_years`, `education_level`
- Модифицируйте валидацию для проверки новых полей
- Добавьте фильтрацию данных по определенным критериям (например, опыт > 3 года)
**Подсказка:**
- Генерация данных происходит в функции `generate_sample_data()` — новые поля добавляются в словарь `data`
- Для `department` подойдёт `random.choice()` со списком отделов, для `experience_years` — привяжите к возрасту (не может быть больше `age - 18`)
- Валидация — в функции `read_and_validate_data()`: добавьте `assert` по аналогии с существующими
- Фильтрацию (pandas `df[df['experience_years'] > 3]`) можно добавить в `transform_data()` и сохранить отдельным файлом
**Цель задания:** Научиться работать с различными типами данных и валидацией.
### Задание 2: Создание дополнительных отчетов
**Сложность:** Средняя
**Задача:**
- Создайте отчет в формате JSON с детальной статистикой по отделам
- Добавьте визуализацию данных с помощью библиотеки matplotlib (сохранение графика в файл)
- Создайте сводку по зарплатам в разных возрастных категориях
**Подсказка:**
- Создайте новую функцию и новый `PythonOperator`, добавьте в цепочку после `summary_task`
- Для JSON-отчёта: загрузите processed_data.csv через pandas, сгруппируйте (`groupby`) и сохраните через `json.dump()`
- matplotlib может быть не установлен в контейнере — это нормально, сделайте эту часть необязательной (обработайте `ImportError`)
- Не забудьте `import pandas as pd` внутри функции (не в шапке файла!)
**Цель задания:** Освоить создание комплексных отчетов и визуализацию.
### Задание 3: Оптимизация обработки
**Сложность:** Продвинутая
**Задача:**
- Разделите обработку данных на параллельные задачи для разных отделов
- Добавьте контроль качества данных (проверка на дубликаты, аномалии)
- Реализуйте механизм повторной обработки при ошибках
**Подсказка:**
- Параллельные задачи задаются списком: `read_task >> [task_a, task_b, task_c] >> summary_task`
- Каждая задача фильтрует DataFrame по своему отделу и сохраняет результат в отдельный файл
- Проверку дубликатов можно сделать через `df['id'].is_unique`
- Параметры `retries` и `retry_delay` можно задать как у отдельных задач, так и в `default_args`
**Цель задания:** Научиться оптимизировать и делать обработку данных отказоустойчивой.
---
## Задания для csv_to_postgres.py
**Цель:** Освоить паттерны загрузки данных (ETL) из файлов в базу данных и применение XCom.
### Задание 1: Расширение структуры данных
**Сложность:** Начальная
**Задача:**
- Добавьте в генератор CSV новую колонку `status` (например, со случайными значениями 'NEW', 'PROCESSING', 'COMPLETED').
- Обновите функцию `_create_table`, чтобы учесть новую колонку.
- Запустите DAG и проверьте, что данные успешно загрузились с новой колонкой.
- Убедитесь, что сгенерированный файл появился в каталоге `/opt/airflow/data/output/`.
**Подсказка:**
- Новая колонка добавляется в трёх местах: генерация (`_generate_csv`), DDL (`_create_table`), загрузка (`_load_csv` — список колонок в COPY)
- В `_generate_csv()` для случайных значений подойдёт `random.choices(['NEW', 'PROCESSING', 'COMPLETED'], k=rows)`
- **Подводный камень:** если таблица уже существует со старой схемой, новая колонка не появится. Удалите таблицу вручную или добавьте `DROP TABLE` перед `CREATE TABLE`
**Цель задания:** Понять процесс изменения схемы данных на всех этапах пайплайна.
### Задание 2: Использование PostgresOperator
**Сложность:** Средняя
**Задача:**
- Перепишите задачу `create_orders_table`. Сейчас она использует `PythonOperator` и `PostgresHook` внутри Python-функции.
- Замените её на использование стандартного `PostgresOperator`, используя соединение `postgres_training`.
- Убедитесь, что пайплайн продолжает работать корректно.
**Подсказка:**
- Посмотрите, как устроен `sql_basic_dag.py` — там `PostgresOperator` уже используется, можно взять за образец
- `PostgresOperator` принимает `postgres_conn_id` и `sql` — SQL можно передать прямо строкой
- После замены проверьте, не осталось ли ссылок на удалённую функцию. `_get_conn()` всё ещё нужна в `_load_csv()`
**Цель задания:** Научиться использовать специализированные операторы для работы с БД вместо кастомного Python-кода.
---
## Задания для csv_to_postgres_dq.py
**Цель:** Освоить подходы к обеспечению качества данных (Data Quality) в Airflow.
### Задание 1: Новая проверка качества
**Сложность:** Средняя
**Задача:**
- Добавьте новую функцию проверки прямо в `csv_to_postgres_dq.py`, которая будет убеждаться, что все значения в колонке `amount` строго больше нуля.
- Добавьте вызов этой функции как новую задачу в DAG `csv_to_postgres_dq`.
- Встройте новую задачу в общую цепочку выполнения (например, перед `data_quality_summary`).
**Подсказка:**
- Возьмите за образец любую из существующих функций проверки (например, `_check_has_rows`) — паттерн одинаковый: подключиться, выполнить SQL, проверить результат, бросить `ValueError` если не ок
- SQL для проверки: `SELECT COUNT(*) FROM public.orders WHERE amount <= 0`
- Не забудьте создать `PythonOperator` и добавить задачу в цепочку
**Цель задания:** Научиться расширять набор проверок качества данных.
### Задание 2: Управление статусом при ошибках (Trigger Rules)
**Сложность:** Продвинутая
**Задача:**
- Смоделируйте ошибку (например, временно измените данные так, чтобы проверки не прошли).
- По умолчанию, если падает одна проверка, следующие не выполняются (поведение `all_success`).
- Измените параметры задач так (с помощью `trigger_rule`), чтобы выполнялись *все* проверки, даже если некоторые из них упали.
- Сделайте так, чтобы задача `data_quality_summary` могла анализировать статусы предыдущих задач и отражать общий итог.
**Подсказка:**
- Для моделирования ошибки: измените `EXPECTED_ORDERS_SCHEMA` так, чтобы схема не совпала. Запустите — увидите, что все задачи после упавшей будут skipped
- Ключевой параметр — `trigger_rule` у `PythonOperator`. Значение `"all_done"` означает: «запустись в любом случае, когда все upstream завершились (успешно или нет)»
- Для анализа статусов в `dq_summary`: функция может принимать `**context` и через `context['ti']` получить информацию о статусах предыдущих задач
- После экспериментов не забудьте вернуть `EXPECTED_ORDERS_SCHEMA` к правильным значениям
**Цель задания:** Освоить продвинутую маршрутизацию статусов задач с помощью `trigger_rule`.
---
## Задания для data_processing_dag.py
**Цель:** Освоить ETL процессы и работу с бизнес-логикой.
### Задание 1: Расширение ETL пайплайна
**Сложность:** Средняя
**Задача:**
- Добавьте новый источник данных — файл с информацией о продуктах
- Создайте задачу для объединения данных о заказах с информацией о продуктах
- Добавьте расчет общей выручки по продуктам
**Подсказка:**
- В `create_sample_data()` добавьте генерацию ещё одного CSV (products.csv) с колонками `product`, `category`, `weight_kg` — значения product должны совпадать с теми, что уже есть в orders
- Создайте функцию `extract_products()` по аналогии с `extract_customers()`
- Для объединения в `transform_data()` используйте `pd.merge()` по ключу `product`
- Новую extract-задачу поставьте параллельно с существующими: `create_data_task >> [extract_customers_task, extract_orders_task, extract_products_task]`
**Цель задания:** Научиться работать с множественными источниками данных.
### Задание 2: Расширение текстового отчёта
**Сложность:** Средняя
**Задача:**
- Расширьте функцию `generate_report()` так, чтобы отчёт содержал больше полезной статистики
- Добавьте в отчёт: топ-3 самых дорогих заказа, распределение заказов по месяцам, среднюю сумму заказа по каждому клиенту
- Сохраните расширенный отчёт в `/opt/airflow/data/output/detailed_report.txt`
**Подсказка:**
- Всё делается в существующей функции `generate_report()` — дополните её
- Полезные методы pandas: `df.nlargest()`, `df.groupby(...).mean()`, `df.groupby(...).sum()`
- Можно записать результат в тот же файл или создать отдельный
**Цель задания:** Освоить создание аналитических отчётов с помощью pandas.
### Задание 3: Мониторинг качества данных
**Сложность:** Продвинутая
**Задача:**
- Добавьте проверки качества данных на каждом этапе ETL
- Создайте механизм оповещения о проблемах с данными
- Реализуйте архивирование обработанных данных
**Подсказка:**
- Создайте отдельную функцию валидации (проверка: строки > 0, нет NULL в ключевых полях, суммы положительные) и поставьте её между extract и transform
- Для оповещения используйте `on_failure_callback` в `default_args` — это функция, которая вызывается при падении любой задачи. Она получает `context` с информацией об ошибке
- Для архивирования: скопируйте результат в файл с датой в имени (например, `shutil.copy()` + `datetime.now().strftime(...)`)
**Цель задания:** Научиться обеспечивать качество данных в ETL процессах.
---
## Задания для branching_dag.py
**Цель:** Освоить условную логику и ветвление в Airflow.
### Задание 1: Модификация условий ветвления
**Сложность:** Средняя
**Задача:**
- Измените условие в функции `check_data_quality()` на основе реальных критериев (например, размер файла)
- Добавьте третью ветку обработки для данных "требующих ручной проверки"
- Настройте разные триггерные правила для слияния веток
**Подсказка:**
- Сейчас функция выбирает ветку случайно — замените `random.random()` на проверку чего-то реального (например, `os.path.getsize()` для размера файла)
- Функция `BranchPythonOperator` возвращает `task_id` ветки для выполнения. Для третьей ветки — верните третий `task_id`
- У `merge_task` проверьте `trigger_rule` — при ветвлении непройденные ветки получают статус `skipped`
**Цель задания:** Научиться создавать сложные условия ветвления.
### Задание 2: Реализация реального сценария
**Сложность:** Продвинутая
**Задача:**
- Создайте реальные задачи обработки для CSV и JSON форматов
- Добавьте валидацию данных в каждой ветке
- Реализуйте механизм сравнения результатов из разных веток
**Подсказка:**
- Вместо `print("Обработка CSV...")` загрузите реальный файл (например, `sample_data.csv`) и посчитайте статистику
- Результат каждой ветки сохраните в отдельный файл — в `merge_results()` загрузите оба (если существуют) и сравните
- Учтите, что при ветвлении выполняется только одна ветка — функция слияния должна обрабатывать случай, когда один из файлов отсутствует
**Цель задания:** Применить ветвление в реальном сценарии обработки данных.
### Задание 3: Динамическое ветвление
**Сложность:** Продвинутая
**Задача:**
- Реализуйте ветвление на основе внешних параметров (например, переданных через Variables)
- Добавьте обработку случая, когда ни одна ветка не подходит
- Создайте механизм логирования выбранного пути выполнения
**Подсказка:**
- Variables создаются в UI: Admin → Variables. Для чтения в коде: `Variable.get('key', default_var='значение')`
- Для случая «ни одна ветка не подходит» добавьте `DummyOperator` как ветку-заглушку
- `BranchPythonOperator` должен всегда возвращать существующий `task_id` — иначе DAG упадёт с ошибкой
**Цель задания:** Освоить динамическое принятие решений в DAG'ах.
---
## Задания для error_handling_dag.py
**Цель:** Освоить обработку ошибок и создание отказоустойчивых пайплайнов.
### Задание 1: Настройка стратегий повторения
**Сложность:** Средняя
**Задача:**
- Измените параметры `retries` и `retry_delay` для разных задач
- Добавьте экспоненциальную задержку между повторными попытками
- Реализуйте кастомный обработчик ошибок для конкретных исключений
**Подсказка:**
- Параметры повторения можно задать на уровне задачи (перекроют `default_args`): `retries`, `retry_delay`, `retry_exponential_backoff`, `max_retry_delay`
- Для кастомного обработчика: параметр `on_retry_callback` принимает функцию с аргументом `context`, из которого можно получить номер попытки через `context['ti'].try_number`
- Запустите DAG несколько раз и посмотрите в логах, как меняется задержка между попытками
**Цель задания:** Научиться настраивать стратегии обработки ошибок.
### Задание 2: Создание комплексной обработки ошибок
**Сложность:** Продвинутая
**Задача:**
- Добавьте задачи для разных типов ошибок (сетевая ошибка, ошибка данных, системная ошибка)
- Создайте механизм эскалации ошибок (после N неудачных попыток)
- Реализуйте отправку уведомлений о критических ошибках
**Подсказка:**
- Создайте функции по аналогии с `unreliable_task()`, но бросающие разные типы исключений: `ConnectionError`, `ValueError`, `OSError`
- Для эскалации используйте `on_failure_callback` — внутри проверьте `context['ti'].try_number` и при достижении порога выведите сообщение уровня CRITICAL
- Пример структуры callback:
```python
def escalation_callback(context):
if context['ti'].try_number >= 3:
# сюда — логику эскалации
```
**Цель задания:** Освоить создание комплексной системы обработки ошибок.
### Задание 3: Добавление логирования через модуль logging
**Сложность:** Средняя
**Задача:**
- Замените все `print()` в функциях DAG'а на вызовы стандартного модуля `logging`
- Добавьте в каждую функцию логирование начала и окончания выполнения
- Посмотрите, как логи отображаются в интерфейсе Airflow (вкладка Log у каждой задачи)
**Подсказка:**
- Стандартный паттерн: `import logging` в начале файла, затем `log = logging.getLogger(__name__)` в функции
- Уровни: `log.info()` для штатных событий, `log.warning()` для предупреждений, `log.error()` для ошибок
- В UI Airflow откройте выполненную задачу → вкладка **Log** — сообщения `logging` отображаются с метками уровня и временем, в отличие от `print()`
**Цель задания:** Научиться использовать стандартное логирование Python в задачах Airflow и читать логи через UI.
---
## Дополнительные задания: пулы, XCom, TaskGroup и алертинг
**Цель:** Освоить продвинутые возможности оркестрации — управление ресурсами, обмен данными между задачами, группировку и оповещения.
### Задание 1: Пулы и управление ресурсами
**Сложность:** Средняя
**Задача:**
- В интерфейсе 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 задач из этого пула
**Подсказка:**
- Пул создаётся в UI: Admin → Pools → "+"
- У `PythonOperator` есть параметры `pool` (имя пула) и `pool_slots` (сколько слотов занимает задача)
- Если задача занимает 2 слота из 2 доступных — вторая задача будет ждать, даже если они не связаны зависимостями
**Цель задания:** Научиться управлять параллелизмом задач через пулы и `pool_slots`.
### Задание 2: Обмен данными через XCom
**Сложность:** Средняя
**Задача:**
- В том же DAG добавьте задачу `calculate_metrics` (PythonOperator), которая возвращает словарь с агрегированными показателями, например: `{"total_orders": ..., "avg_amount": ...}`
- Добавьте задачу `log_metrics`, которая с помощью `xcom_pull` читает результат `calculate_metrics` и выводит значения в лог
- Для одной из задач продемонстрируйте использование XCom в Jinja-шаблоне (например, в `bash_command` или SQL-запросе)
**Подсказка:**
- Любое значение, которое функция возвращает через `return`, автоматически сохраняется в XCom
- Для чтения XCom в другой задаче: функция принимает `**context`, затем `context['ti'].xcom_pull(task_ids='имя_задачи')`
- В Jinja-шаблонах (параметры `bash_command`, `sql` и др.) XCom доступен через `{{ ti.xcom_pull(task_ids='...') }}`
- Посмотреть сохранённые XCom-значения можно в UI: откройте задачу → вкладка **XCom**
**Цель задания:** Освоить передачу результатов между задачами через XCom и их использование в шаблонах.
### Задание 3: TaskGroup и алертинг
**Сложность:** Продвинутая
**Задача:**
- Объедините логически связанные задачи (например, `extract` / `transform` / `load`) в `TaskGroup`'ы
- Добавьте завершающую задачу `send_notification` на основе примера из раздела про алертинг (PythonOperator с `send_email_smtp` или другим механизмом уведомлений)
- Настройте для задачи уведомления `trigger_rule=TriggerRule.ALL_DONE`, чтобы уведомление отправлялось даже при частичных ошибках
- При желании вынесите функцию отправки письма в отдельный модуль `utils.py` и импортируйте её в DAG
**Подсказка:**
- TaskGroup — контекстный менеджер: `with TaskGroup("имя") as group:` — внутри определяете задачи как обычно
- Зависимости работают на уровне групп: `group_a >> group_b`
- `TriggerRule.ALL_DONE` (из `airflow.utils.trigger_rule`) означает: «запустись, когда все upstream завершились — неважно, с успехом или ошибкой»
- В Graph View группы отображаются как складные блоки — удобно для больших DAG'ов
**Цель задания:** Научиться группировать задачи с помощью TaskGroup и строить схему оповещений о статусе пайплайна.
---
## Рекомендации по выполнению
1. **Начинайте с простых заданий** и постепенно переходите к сложным
2. **Тестируйте каждое изменение** через интерфейс Airflow
3. **Изучайте логи выполнения** для понимания поведения задач (вкладка Log у каждой задачи в UI)
4. **Экспериментируйте** с разными настройками и параметрами
5. **Не бойтесь ломать** — DAG'и можно откатить через `git checkout`, а таблицы пересоздать
## Оценка прогресса
- **Начальный уровень:** Выполнены задания для hello_world_dag и sql_basic_dag
- **Средний уровень:** Выполнены задания для file_operations_dag, csv_to_postgres.py, csv_to_postgres_dq.py и data_processing_dag
- **Продвинутый уровень:** Выполнены все задания, включая branching_dag, error_handling_dag и дополнительные задания по пулам, XCom, TaskGroup и алертингу
Удачи в изучении Apache Airflow!