diff --git a/09 - Продвинутые возможности Airflow.md b/09 - Продвинутые возможности Airflow.md index f2df916..f1ad8ea 100644 --- a/09 - Продвинутые возможности Airflow.md +++ b/09 - Продвинутые возможности Airflow.md @@ -1,6 +1,6 @@ -# Продвинутые возможности Airflow для начинающих специалистов +# Пулы, XCom, TaskGroup и алертинг в Airflow -В этом материале мы познакомимся с расширенными функциями Airflow, которые помогут вам решать нетривиальные задачи при работе со сложными пайплайнами. +В этом материале мы познакомимся с управлением ресурсами через пулы задач, обменом данными между задачами через XCom, группировкой задач с помощью TaskGroup и настройкой системы оповещений (алертинга) в Airflow. Начнем с улучшения нашего базового пайплайна — проведем рефакторинг кода для лучшей читаемости и поддержки. @@ -12,11 +12,18 @@ Когда планировщик обнаруживает, что наступило время выполнения DAG, он запускает задачу согласно заданной последовательности. При этом задача занимает один слот в пуле и освобождает его после завершения. -Создать новый пул задач и установить его размер можно через веб-интерфейс Airflow: +Создать новый пул задач и установить его размер можно через веб-интерфейс Airflow. Рассмотрим пример пула `data_processing_pool` для тяжёлых задач: -![image](_attachments/advanced_feature_demo_1.gif) +1. В верхнем меню откройте **Admin → Pools**. +2. Нажмите кнопку **+** или **Create**. +3. В поле **Name** укажите, например, `data_processing_pool`. +4. В поле **Slots** задайте количество слотов, например `5` + (это значит, что одновременно смогут выполняться не более 5 задач из этого пула). +5. При желании заполните **Description** — например, `Пул для тяжёлых задач бэкапа`. +6. Нажмите **Save**. + +После этого в задачах можно указать параметр `pool="data_processing_pool"`, и они будут занимать слоты именно этого пула. -*Создание нового пула задач в Airflow* Зачем нужны пулы задач? Они помогают: - Организовать запуск процессов в системе @@ -59,13 +66,34 @@ BashOperator( 💡 **Контекст выполнения** — это набор параметров, передаваемых при запуске задачи, а также метаданные, генерируемые во время выполнения: время старта, время завершения, имя DAG и другие. Контекст можно найти следующим образом: -![image](_attachments/advanced_feature_demo_3.gif){:height 437, :width 778} +1. Откройте веб-интерфейс Airflow и перейдите на вкладку **DAGs**. +2. Найдите нужный DAG и кликните по его имени. +3. На странице DAG убедитесь, что открыта вкладка **Grid**. +4. В сетке выберите нужный запуск DAG (колонка) и задачу (строка) и кликните по цветному квадратику задачи. +5. Справа откроется панель **Task Instance** (детали экземпляра задачи). В ней можно увидеть: + - `dag_id` и `task_id`; + - логическую дату запуска (**logical_date / execution_date**); + - текущий статус, время старта и завершения; + - ссылки на лог, XCom и другую служебную информацию. -В XCom можно передавать сериализованные объекты. Значения XCom хранятся в базе данных Airflow и доступны через интерфейс: +Именно эти поля и составляют большую часть «контекста выполнения», который доступен в Jinja-шаблонах через объекты вроде `{{ dag_run }}` и `{{ task_instance }}`. -![image](_attachments/advanced_feature_demo_2.gif) +В XCom можно передавать сериализованные объекты. Значения XCom хранятся в базе данных Airflow и доступны через интерфейс. -*XCom в веб-интерфейсе Airflow* +Посмотреть XCom можно двумя способами. + +**1. Через конкретную задачу в DAG** + +1. Откройте нужный DAG и вкладку **Grid**. +2. Найдите нужный запуск и кликните по квадратику задачи. +3. В правой панели **Task Instance** перейдите на вкладку **XCom**. +4. В таблице вы увидите все XCom-записи для этого экземпляра задачи: ключ (`key`), значение (`value`), время создания и т. д. + +**2. Через общий список XCom** + +1. В верхнем меню выберите **Browse → XComs**. +2. Отфильтруйте записи по `dag_id`, `task_id` или другим полям, если нужно. +3. Откройте интересующую запись, чтобы увидеть её содержимое. **Важное правило**: XCom предназначен для обмена небольшими сообщениями. Данные проходят сериализацию/десериализацию при чтении и записи в таблицу. @@ -73,7 +101,7 @@ BashOperator( Для передачи больших объемов данных используйте внешние средства: файловую систему, базы данных (чаще всего PostgreSQL) или очереди сообщений (например, Kafka). -Механизм XCom похож на работу функций в Python. Многие операторы по умолчанию возвращают результат выполнения задачи. За это отвечает параметр `do_xcom_push`, который по умолчанию равен `True`. +Механизм XCom похож на работу функций в Python. Многие операторы (например, `PythonOperator`) по умолчанию возвращают результат выполнения задачи. За это отвечает параметр `do_xcom_push`, который во многих случаях равен `True` по умолчанию. Чтобы прочитать сообщения из XCom, используйте метод `xcom_pull` в контексте задачи: @@ -149,17 +177,7 @@ with TaskGroup("data_loading") as loading_group: Обратите внимание: при использовании TaskGroup последовательность задач указывается внутри группы после объявления всех задач, а в конце DAG описывается последовательность выполнения самих групп. -Визуально в интерфейсе Airflow это выглядит так: - -![image](_attachments/advanced_feature_diagram_17.png) - -*Группировка задач с помощью TaskGroup* - -TaskGroup добавляет интерактивность в веб-интерфейс — группы задач можно сворачивать и разворачивать, что значительно улучшает восприятие DAG с большим количеством задач и связей: - -![image](_attachments/advanced_feature_animation_23.gif) - -*Группу задач можно раскрыть для детального просмотра* +Визуально в интерфейсе Airflow группы задач отображаются как один узел с небольшим индикатором. Клик по нему разворачивает или сворачивает вложенные задачи, что значительно улучшает восприятие DAG с большим количеством задач и связей, особенно когда в них десятки и сотни задач. TaskGroup — это удобный способ логической группировки задач, который помогает упростить код и представить сложные пайплайны более компактно. @@ -167,7 +185,7 @@ TaskGroup — это удобный способ логической групп Алертинг — один из ключевых компонентов системы оркестрации, так как важно своевременно получать уведомления об ошибках для их оперативного анализа и решения. -По умолчанию в Airflow настроена отправка уведомлений на электронную почту. При создании DAG указываются email-адреса, на которые будут отправляться сообщения. С помощью параметров можно настроить различные сценарии оповещений. +В Airflow есть встроенная поддержка отправки уведомлений на электронную почту (при условии, что в конфигурации настроен SMTP-сервер). При создании DAG указываются email-адреса, на которые будут отправляться сообщения. С помощью параметров можно настроить различные сценарии оповещений. Давайте модифицируем наш первый DAG так, чтобы получать уведомления на почту при возникновении ошибок. При этом настроим перезапуск задач в случае неудачи (например, 2 попытки), но без уведомлений о самих перезапусках: @@ -208,14 +226,16 @@ dag = DAG( from datetime import datetime, timedelta, timezone import dateutil from airflow.utils.email import send_email_smtp +from airflow.operators.python import get_current_context MAIL_LIST = [ "email_1@gmail.ru", "email_2@gmail.ru" ] -def notify_email(calculation_dt: str, dagrun_begin_time, **context): +def notify_email(calculation_dt: str, dagrun_begin_time: str): # dag_run date is utc timezone, so add `timezone.utc` to calculation_end to combat 3 hour diff. + context = get_current_context() calculation_start = dateutil.parser.isoparse(dagrun_begin_time.replace('Z', '+00:00')) calculation_end = datetime.now(timezone.utc) duration = str(timedelta(seconds=(calculation_end - calculation_start).seconds)) @@ -241,6 +261,7 @@ def notify_email(calculation_dt: str, dagrun_begin_time, **context): ```python from airflow.operators.python import PythonOperator +from airflow.utils.trigger_rule import TriggerRule from utils import notify_email LOCAL_CALCULATION_DT = '{{ dag.timezone.convert(execution_date).strftime("%Y-%m-%d") }}' @@ -249,7 +270,6 @@ DAG_RUN_BEGIN_TIME = "{{ dag_run.start_date }}" email_notification_task = PythonOperator( task_id="send_email_notification", python_callable=notify_email, - provide_context=True, dag=dag, trigger_rule=TriggerRule.ALL_DONE, op_args=[LOCAL_CALCULATION_DT, DAG_RUN_BEGIN_TIME], @@ -267,6 +287,7 @@ email_notification_task = PythonOperator( В этом примере используется переменная 'database_connection_string', предварительно созданная в интерфейсе Airflow, для хранения строки подключения к базе данных. Это позволяет избежать жесткого кодирования конфиденциальной информации в коде DAG и упрощает настройку подключения для разных окружений. ```python +import os import datetime as dt import pandas as pd from airflow.models import DAG @@ -351,6 +372,7 @@ start_task >> data_processing Итоговый DAG будет выглядеть следующим образом: ```python +import os import datetime as dt import pandas as pd from airflow.models import DAG @@ -426,4 +448,4 @@ start_task >> data_processing - Логическая группировка задач с TaskGroup - Настройка системы оповещений -Помните: не стоит использовать все доступные функции сразу. Выбирайте инструменты последовательно и находите оптимальный набор возможностей под конкретную задачу. \ No newline at end of file +Помните: не стоит использовать все доступные функции сразу. Выбирайте инструменты последовательно и находите оптимальный набор возможностей под конкретную задачу. diff --git a/_attachments/advanced_feature_animation_23.gif b/_attachments/advanced_feature_animation_23.gif deleted file mode 100644 index 8ae977f..0000000 Binary files a/_attachments/advanced_feature_animation_23.gif and /dev/null differ diff --git a/_attachments/advanced_feature_demo_1.gif b/_attachments/advanced_feature_demo_1.gif deleted file mode 100644 index 62f991d..0000000 Binary files a/_attachments/advanced_feature_demo_1.gif and /dev/null differ diff --git a/_attachments/advanced_feature_demo_2.gif b/_attachments/advanced_feature_demo_2.gif deleted file mode 100644 index 2e205bf..0000000 Binary files a/_attachments/advanced_feature_demo_2.gif and /dev/null differ diff --git a/_attachments/advanced_feature_demo_3.gif b/_attachments/advanced_feature_demo_3.gif deleted file mode 100644 index 672762f..0000000 Binary files a/_attachments/advanced_feature_demo_3.gif and /dev/null differ diff --git a/_attachments/advanced_feature_diagram_17.png b/_attachments/advanced_feature_diagram_17.png deleted file mode 100644 index 6cee324..0000000 Binary files a/_attachments/advanced_feature_diagram_17.png and /dev/null differ