- Зачем: - упражнения должны быть понятнее и соответствовать ожидаемому уровню сложности. - Что: - добавлены подсказки и описания подводных камней. - упрощены задания по отчетам, мониторингу и логированию. - Проверка: - git diff --check.
36 KiB
Учебные задания по Apache Airflow
Этот документ содержит практические задания для закрепления знаний по Apache Airflow. Каждое задание соответствует одному из учебных DAG'ов и направлено на лучшее понимание конкретных концепций.
Общие инструкции
Перед выполнением заданий:
- Убедитесь, что стенд Airflow запущен:
docker-compose up -d - Проверьте доступность интерфейса: http://localhost:8080
- Ознакомьтесь с соответствующим DAG'ом в интерфейсе Airflow
- Откройте исходный код 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 и строить схему оповещений о статусе пайплайна.
Рекомендации по выполнению
- Начинайте с простых заданий и постепенно переходите к сложным
- Тестируйте каждое изменение через интерфейс Airflow
- Изучайте логи выполнения для понимания поведения задач (вкладка Log у каждой задачи в UI)
- Экспериментируйте с разными настройками и параметрами
- Не бойтесь ломать — 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!