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

36 KiB
Raw Permalink Blame History

Учебные задания по 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:
    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!