From 97241b5b035793b71063abb17896662c63176f50 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sat, 6 Dec 2025 21:36:21 +0300 Subject: [PATCH] =?UTF-8?q?=D0=9F=D0=B5=D1=80=D0=B5=D0=B8=D0=BC=D0=B5?= =?UTF-8?q?=D0=BD=D0=BE=D0=B2=D0=B0=D0=BD=D0=B8=D0=B5=20=D1=80=D0=B0=D0=B7?= =?UTF-8?q?=D0=B4=D0=B5=D0=BB=D0=B0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- ...ы, XCom, TaskGroup и алертинг в Airflow.md | 488 ++++++++++++++++++ 1 file changed, 488 insertions(+) create mode 100644 09 - Пулы, XCom, TaskGroup и алертинг в Airflow.md diff --git a/09 - Пулы, XCom, TaskGroup и алертинг в Airflow.md b/09 - Пулы, XCom, TaskGroup и алертинг в Airflow.md new file mode 100644 index 0000000..8e1443b --- /dev/null +++ b/09 - Пулы, XCom, TaskGroup и алертинг в Airflow.md @@ -0,0 +1,488 @@ +# Пулы, XCom, TaskGroup и алертинг в Airflow + +В этом материале мы познакомимся с управлением ресурсами через пулы задач, обменом данными между задачами через XCom, группировкой задач с помощью TaskGroup и настройкой системы оповещений (алертинга) в Airflow. + +Начнем с улучшения нашего базового пайплайна — проведем рефакторинг кода для лучшей читаемости и поддержки. + +## Управление ресурсами с помощью пулов задач + +В системах с высокой нагрузкой, где одновременно запускается множество задач и DAG-ов, может возникнуть чрезмерная нагрузка на исполнителей и серверную часть. Это может привести к ошибкам выполнения и даже к отказу системы, если не установить соответствующие ограничения. + +В Airflow для решения этой проблемы существует механизм управления ресурсами — **пулы задач** (pools). По умолчанию в Airflow настроен один пул задач — `default_pool` с 128 слотами, что означает возможность параллельного выполнения 128 задач одновременно. Пул `default_pool` нельзя удалить, но можно изменить его размер — увеличить или уменьшить количество слотов. + +Когда планировщик обнаруживает, что наступило время выполнения DAG, он запускает задачу согласно заданной последовательности. При этом задача занимает один слот в пуле и освобождает его после завершения. + +Создать новый пул задач и установить его размер можно через веб-интерфейс Airflow. Рассмотрим пример пула `data_processing_pool` для тяжёлых задач: + +1. В верхнем меню откройте **Admin → Pools**. +2. Нажмите кнопку **+** или **Create**. +3. В поле **Name** укажите, например, `data_processing_pool`. +4. В поле **Slots** задайте количество слотов, например `5` + (это значит, что одновременно смогут выполняться не более 5 задач из этого пула). +5. При желании заполните **Description** — например, `Пул для тяжёлых задач бэкапа`. +6. Нажмите **Save**. + +После этого в задачах можно указать параметр `pool="data_processing_pool"`, и они будут занимать слоты именно этого пула. + + +Зачем нужны пулы задач? Они помогают: +- Организовать запуск процессов в системе +- Предотвратить перегрузку системы при выполнении большого количества ресурсоемких задач + +Вы можете задать "вес" задачи через параметр `pool_slots`, чтобы оптимизировать распределение нагрузки. Если общее количество задач превышает доступные слоты, планировщик поставит задачу в очередь и запустит ее, как только появятся свободные ресурсы. + +Пример настройки веса задач: + +В приведенном примере показано, как можно настроить использование пула задач с различным весом. Задача 'data_backup_task' использует 3 слота в пуле 'data_processing_pool', что делает ее более ресурсоемкой по сравнению с задачами 'file_check_task' и 'cleanup_task', которые используют по 1 слоту. Это позволяет контролировать распределение ресурсов между различными задачами. + +```python +BashOperator( + task_id="data_backup_task", + bash_command="bash backup_script.sh", + pool_slots=3, + pool="data_processing_pool", +) + +BashOperator( + task_id="file_check_task", + bash_command="bash validate_files.sh", + pool_slots=1, + pool="data_processing_pool", +) + +BashOperator( + task_id="cleanup_task", + bash_command="bash cleanup_files.sh", + pool_slots=1, + pool="data_processing_pool", +) +``` + +Более подробную информацию о механизме пулов можно найти в [официальной документации](https://airflow.apache.org/docs/apache-airflow/stable/concepts/pools.html#pools). + +## Обмен данными между задачами: XCom и контекст выполнения + +**XCom** (от англ. cross-communications — «межзадачная коммуникация») — это механизм обмена сообщениями между задачами внутри одного DAG. Архитектурно каждая задача изолирована от других и работает в собственном контексте. + +💡 **Контекст выполнения** — это набор параметров, передаваемых при запуске задачи, а также метаданные, генерируемые во время выполнения: время старта, время завершения, имя DAG и другие. Контекст можно найти следующим образом: + +1. Откройте веб-интерфейс Airflow и перейдите на вкладку **DAGs**. +2. Найдите нужный DAG и кликните по его имени. +3. На странице DAG убедитесь, что открыта вкладка **Grid**. +4. В сетке выберите нужный запуск DAG (колонка) и задачу (строка) и кликните по цветному квадратику задачи. +5. Справа откроется панель **Task Instance** (детали экземпляра задачи). В ней можно увидеть: + - `dag_id` и `task_id`; + - логическую дату запуска (**logical_date / execution_date**); + - текущий статус, время старта и завершения; + - ссылки на лог, XCom и другую служебную информацию. + +Именно эти поля и составляют большую часть «контекста выполнения», который доступен в Jinja-шаблонах через объекты вроде `{{ dag_run }}` и `{{ task_instance }}`. + +В XCom можно передавать сериализованные объекты. Значения 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 предназначен для обмена небольшими сообщениями. Данные проходят сериализацию/десериализацию при чтении и записи в таблицу. + +💡 **Сериализация** — процесс преобразования структуры данных в последовательность байтов. **Десериализация** — восстановление структуры данных из байтовой последовательности. + +Для передачи больших объемов данных используйте внешние средства: файловую систему, базы данных (чаще всего PostgreSQL) или очереди сообщений (например, Kafka). + +Механизм XCom похож на работу функций в Python. Многие операторы (например, `PythonOperator`) по умолчанию возвращают результат выполнения задачи. За это отвечает параметр `do_xcom_push`, который во многих случаях равен `True` по умолчанию. + +Чтобы прочитать сообщения из XCom, используйте метод `xcom_pull` в контексте задачи: + +В этом примере мы получаем результат выполнения задачи с идентификатором 'data_processing_task' с помощью метода xcom_pull. Это позволяет передавать небольшие объемы данных между задачами в рамках одного DAG. + +```python +result = task_instance.xcom_pull(task_ids='data_processing_task') +``` + +Также можно обращаться к сообщениям через Jinja-шаблоны: + +``` +SELECT * FROM {{ task_instance.xcom_pull(task_ids='foo', key='table_name') }} +``` + +Если оператор возвращает значение и параметр `do_xcom_push` установлен в `True` (по умолчанию), это значение автоматически записывается в XCom. + +Пример явного запрета записи в XCom: + +В этом примере задача 'show_directory_contents' создается с параметром do_xcom_push=False, что означает, что результат выполнения этой задачи не будет автоматически сохранен в XCom. Это полезно, когда вы не хотите, чтобы задача передавала какие-либо данные другим задачам через XCom. + +```python +list_files = BashOperator( + task_id='show_directory_contents', + bash_command='ls -la', + do_xcom_push=False +) + +def calculate_sum(): + return 2 + 3 +``` + +Пример автоматической записи в XCom (параметр `do_xcom_push` по умолчанию `True`): + +В этом примере задача 'calculate_sum' автоматически записывает результат выполнения функции calculate_sum в XCom, так как параметр do_xcom_push по умолчанию установлен в True. Это позволяет использовать результат этой задачи в других задачах DAG через XCom. + +```python +sum_result = PythonOperator( + task_id='calculate_sum', + python_callable=calculate_sum, +) +``` + +XCom похожи на переменные (variables) в Airflow, но предназначены именно для взаимодействия между задачами в рамках одного DAG, а не для глобальных настроек. + +XCom упрощает взаимодействие между задачами и применяется в различных сценариях. + +## Группировка задач с помощью TaskGroup + +Для удобства визуализации задачи можно группировать в веб-интерфейсе Airflow (начиная с версии 2.0). Повторяющиеся или логически связанные задачи можно объединить в группы: + +В этом примере задачи сгруппированы в три логические группы: 'data_extraction', 'data_transformation' и 'data_loading'. Каждая группа содержит несколько задач, которые выполняются последовательно внутри группы. Затем группы связаны между собой, чтобы показать общий порядок выполнения этапов обработки данных. + +```python +with TaskGroup("data_extraction") as extraction_group: + extract_1 = DummyOperator(task_id="extract_source_1") + extract_2 = DummyOperator(task_id="extract_source_2") + extract_3 = DummyOperator(task_id="extract_source_3") + extract_1 >> extract_2 >> extract_3 + +with TaskGroup("data_transformation") as transformation_group: + transform_1 = DummyOperator(task_id="transform_step_1") + transform_2 = DummyOperator(task_id="transform_step_2") + transform_1 >> transform_2 + +with TaskGroup("data_loading") as loading_group: + load_1 = DummyOperator(task_id="load_to_target_1") + load_2 = DummyOperator(task_id="load_to_target_2") + load_1 >> load_2 + +[extraction_group, transformation_group] >> loading_group +``` + +Обратите внимание: при использовании TaskGroup последовательность задач указывается внутри группы после объявления всех задач, а в конце DAG описывается последовательность выполнения самих групп. + +Ниже приведена упрощённая схема зависимостей между тремя группами задач: + +```mermaid +flowchart LR + + %% group1 + subgraph G1["group1"] + g1_t1["task1"] + g1_t2["task2"] + g1_t3["task3"] + + g1_t1 --> g1_t2 + g1_t1 --> g1_t3 + end + + %% group2 + subgraph G2["group2"] + g2_t1["task1"] + g2_t2["task2"] + + g2_t1 --> g2_t2 + end + + %% group3 + subgraph G3["group3"] + g3_t1["task1"] + g3_t2["task2"] + + g3_t1 --> g3_t2 + end + + %% зависимости между группами + g1_t2 --> g3_t1 + g1_t3 --> g3_t1 + g2_t2 --> g3_t1 +``` + +Визуально в интерфейсе Airflow группы задач отображаются как один узел с небольшим индикатором. Клик по нему разворачивает или сворачивает вложенные задачи, что значительно улучшает восприятие DAG с большим количеством задач и связей, особенно когда в них десятки и сотни задач. + +TaskGroup — это удобный способ логической группировки задач, который помогает упростить код и представить сложные пайплайны более компактно. + +## Система оповещений (алертинг) + +Алертинг — один из ключевых компонентов системы оркестрации, так как важно своевременно получать уведомления об ошибках для их оперативного анализа и решения. + +В Airflow есть встроенная поддержка отправки уведомлений на электронную почту (при условии, что в конфигурации настроен SMTP-сервер). При создании DAG указываются email-адреса, на которые будут отправляться сообщения. С помощью параметров можно настроить различные сценарии оповещений. + +Давайте модифицируем наш первый DAG так, чтобы получать уведомления на почту при возникновении ошибок. При этом настроим перезапуск задач в случае неудачи (например, 2 попытки), но без уведомлений о самих перезапусках: + +В приведенном примере создан DAG с идентификатором 'customer_analysis_pipeline', который настроен на отправку уведомлений по электронной почте только при ошибках (email_on_failure=True), но не при повторных попытках (email_on_retry=False). Также установлено 2 попытки повторного запуска задач при ошибках с задержкой 2 минуты между попытками. + +```python +import os +import datetime as dt +import pandas as pd +from airflow.models import DAG +from airflow.operators.python import PythonOperator +from airflow.operators.bash import BashOperator +from sqlalchemy import create_engine + +# основные параметры DAG +args = { + 'owner': 'data_engineering_team', + 'start_date': dt.datetime(2021, 6, 15), + 'retries': 2, + 'retry_delay': dt.timedelta(minutes=2), + 'email': ["data-team@example.com"], + 'email_on_failure': True, + 'email_on_retry': False, +} + +dag = DAG( + dag_id='customer_analysis_pipeline', + schedule_interval=None, + default_args=args, +) +``` + +Теперь вы будете получать email-уведомления при ошибках выполнения. + +Пример функции для отправки уведомлений с использованием параметров из контекста: + +```python +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: 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)) + + title = f"Ежедневные расчёты завершились успешно (dag: {context['task_instance'].dag_id})." + + body = f""" + Привет,
+ Я закончил работу над расчётом за "{calculation_dt}". Это заняло {duration} часов/минут/секунд.
+ Логи и запуски тоже можно посмотреть тут. +
+ +
+
+ Навеки твой,
+ Airflow бот
+ """ + + send_email_smtp(";".join(MAIL_LIST), title, body) +``` + +Такую функцию можно разместить в отдельном файле (например, `utils.py`), импортировать как модуль в нужных DAG и вызывать отдельной задачей: + +```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") }}' +DAG_RUN_BEGIN_TIME = "{{ dag_run.start_date }}" + +email_notification_task = PythonOperator( + task_id="send_email_notification", + python_callable=notify_email, + dag=dag, + trigger_rule=TriggerRule.ALL_DONE, + op_args=[LOCAL_CALCULATION_DT, DAG_RUN_BEGIN_TIME], +) + +... >> email_notification_task +``` + +## Практическое применение: улучшенный DAG + +Теперь применим изученные концепции для усовершенствования нашего DAG. Добавим переменные и разобьем задачи на логические группы. + +Сначала создадим переменную `DATABASE_URL` со строкой подключения к базе данных через веб-интерфейс Airflow и импортируем ее в коде DAG: + +В этом примере используется переменная 'database_connection_string', предварительно созданная в интерфейсе Airflow, для хранения строки подключения к базе данных. Это позволяет избежать жесткого кодирования конфиденциальной информации в коде DAG и упрощает настройку подключения для разных окружений. + +```python +import os +import datetime as dt +import pandas as pd +from airflow.models import DAG +from airflow.operators.bash import BashOperator +from airflow.operators.python import PythonOperator +from airflow.operators.dummy import DummyOperator +from airflow.utils.task_group import TaskGroup +from airflow.models import Variable +from sqlalchemy import create_engine + +DATABASE_URL = Variable.get('database_connection_string') + +args = { + 'owner': 'analytics_team', + 'start_date': dt.datetime(2021, 6, 15), + 'retries': 2, + 'retry_delay': dt.timedelta(minutes=2), +} + +# функции для обработки данных +def get_file_path(file_name): + return os.path.join(os.path.expanduser('~/data'), file_name) + +def load_customer_data(): + file_path = get_file_path('customer_data.csv') + df = pd.read_csv(file_path) + engine = create_engine(DATABASE_URL) + df.to_sql('customers', engine, index=False, if_exists='replace', schema='staging') + +def aggregate_customer_data(): + engine = create_engine(DATABASE_URL) + customer_df = pd.read_sql('select * from staging.customers', con=engine) + + df = customer_df.groupby(['region', 'category']).agg( + total_orders=('orders', 'sum'), + avg_amount=('amount', 'mean') + ).reset_index() + + df.to_sql('customer_summary', engine, index=False, if_exists='replace', schema='analytics') + +dag = DAG( + dag_id='customer_pipeline_enhanced', + schedule_interval=None, + default_args=args, +) +``` + +Теперь разобьем задачи на логические группы и добавим Jinja-шаблоны для доступа к контексту выполнения: + +В этом примере создается начальная задача 'pipeline_start', которая использует Jinja-шаблоны для вывода информации о запуске DAG, включая идентификатор запуска (run_id) и информацию о DAG Run. Затем задачи группируются в логическую группу 'data_processing_stage', что улучшает структуру и читаемость DAG. + +```python +# Начальная задача с информацией о запуске +start_task = BashOperator( + task_id='pipeline_start', + bash_command='echo "Pipeline started! Run ID: {{ run_id }} | DAG Run: {{ dag_run }}"', + dag=dag, +) + +# Группа задач по предварительной обработке данных +with TaskGroup(group_id="data_processing_stage") as data_processing: + # Загрузка данных + load_customer_dataset = PythonOperator( + task_id='load_customer_data', + python_callable=load_customer_data, + dag=dag, + ) + # Агрегация и запись данных + aggregate_customer_dataset = PythonOperator( + task_id='aggregate_customer_data', + python_callable=aggregate_customer_data, + dag=dag, + ) + load_customer_dataset >> aggregate_customer_dataset + +# Установка последовательности выполнения +start_task >> data_processing +``` + +Поскольку последовательность задач внутри групп указывается при их создании, в конце необходимо определить порядок выполнения самих групп, чтобы планировщик понимал общую логику выполнения. + +Итоговый DAG будет выглядеть следующим образом: + +```python +import os +import datetime as dt +import pandas as pd +from airflow.models import DAG +from airflow.operators.bash import BashOperator +from airflow.operators.python import PythonOperator +from airflow.operators.dummy import DummyOperator +from airflow.utils.task_group import TaskGroup +from airflow.models import Variable +from sqlalchemy import create_engine + +DATABASE_URL = Variable.get('database_connection_string') + +args = { + 'owner': 'analytics_team', + 'start_date': dt.datetime(2021, 6, 15), + 'retries': 2, + 'retry_delay': dt.timedelta(minutes=2), +} + +def get_file_path(file_name): + return os.path.join(os.path.expanduser('~/data'), file_name) + +def load_customer_data(): + file_path = get_file_path('customer_data.csv') + df = pd.read_csv(file_path) + engine = create_engine(DATABASE_URL) + df.to_sql('customers', engine, index=False, if_exists='replace', schema='staging') + +def aggregate_customer_data(): + engine = create_engine(DATABASE_URL) + customer_df = pd.read_sql('select * from staging.customers', con=engine) + + df = customer_df.groupby(['region', 'category']).agg( + total_orders=('orders', 'sum'), + avg_amount=('amount', 'mean') + ).reset_index() + + df.to_sql('customer_summary', engine, index=False, if_exists='replace', schema='analytics') + +dag = DAG( + dag_id='customer_pipeline_enhanced', + schedule_interval=None, + default_args=args, +) + +# Начальная задача +start_task = BashOperator( + task_id='pipeline_start', + bash_command='echo "Pipeline started! Run ID: {{ run_id }} | DAG Run: {{ dag_run }}"', + dag=dag, +) + +# Группа предварительной обработки +with TaskGroup(group_id="data_processing_stage") as data_processing: + load_customer_dataset = PythonOperator( + task_id='load_customer_data', + python_callable=load_customer_data, + dag=dag, + ) + aggregate_customer_dataset = PythonOperator( + task_id='aggregate_customer_data', + python_callable=aggregate_customer_data, + dag=dag, + ) + load_customer_dataset >> aggregate_customer_dataset + +start_task >> data_processing +``` + +В этом материале мы рассмотрели расширенные возможности Airflow, которые помогут улучшить работу ваших пайплайнов: +- Управление ресурсами с помощью пулов задач +- Обмен данными между задачами через XCom +- Логическая группировка задач с TaskGroup +- Настройка системы оповещений + +Помните: не стоит использовать все доступные функции сразу. Выбирайте инструменты последовательно и находите оптимальный набор возможностей под конкретную задачу.