From eaae2e1beb8f0c7df67ab07d56ae222fe82f1cd0 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sat, 6 Dec 2025 21:36:35 +0300 Subject: [PATCH] =?UTF-8?q?=D0=90=D0=BA=D1=82=D1=83=D0=B0=D0=BB=D0=B8?= =?UTF-8?q?=D0=B7=D0=B0=D1=86=D0=B8=D1=8F=20=D0=B2=D0=B5=D1=80=D1=81=D0=B8?= =?UTF-8?q?=D0=B8=20Ariflow?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- 01 - Введение в Airflow.md | 2 +- 09 - Продвинутые возможности Airflow.md | 488 ------------------------ README.md | 13 +- 3 files changed, 8 insertions(+), 495 deletions(-) delete mode 100644 09 - Продвинутые возможности Airflow.md diff --git a/01 - Введение в Airflow.md b/01 - Введение в Airflow.md index 1a8dcc5..d83cea8 100644 --- a/01 - Введение в Airflow.md +++ b/01 - Введение в Airflow.md @@ -85,4 +85,4 @@ flowchart LR В этом материале вы познакомились с предпосылками появления инструментов управления процессами, подобных Airflow. Вы узнали, в каких сценариях Airflow наиболее эффективен, а в каких случаях стоит рассмотреть альтернативные решения, а также как этот инструмент применяется на практике в крупных российских и международных компаниях. -💡 Обратите внимание, что в данном модуле рассматривается Airflow версии 2.5. \ No newline at end of file +💡 Обратите внимание, что в данном модуле рассматривается Airflow версии 2.9. diff --git a/09 - Продвинутые возможности Airflow.md b/09 - Продвинутые возможности Airflow.md deleted file mode 100644 index 8e1443b..0000000 --- a/09 - Продвинутые возможности Airflow.md +++ /dev/null @@ -1,488 +0,0 @@ -# Пулы, 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 -- Настройка системы оповещений - -Помните: не стоит использовать все доступные функции сразу. Выбирайте инструменты последовательно и находите оптимальный набор возможностей под конкретную задачу. diff --git a/README.md b/README.md index 0f70cba..66660a9 100644 --- a/README.md +++ b/README.md @@ -11,7 +11,7 @@ Apache Airflow — это мощный оркестратор рабочих п - Хотите освоить свой первый ETL-инструмент - Нуждаетесь в практическом руководстве с понятными объяснениями -Мы начнем с основ и постепенно перейдем к продвинутым возможностям, чтобы вы могли уверенно использовать Airflow в реальных проектах. Курс ориентирован на версию Airflow 2.5 и использует практический подход с множеством примеров и визуальных материалов. +Мы начнем с основ и постепенно перейдем к продвинутым возможностям, чтобы вы могли уверенно использовать Airflow в реальных проектах. Курс ориентирован на версию Airflow 2.9 и использует практический подход с множеством примеров и визуальных материалов. ## 🎯 Что вы узнаете @@ -21,7 +21,7 @@ Apache Airflow — это мощный оркестратор рабочих п - **Интерфейс**: как эффективно использовать веб-интерфейс для мониторинга и управления - **Практическое создание DAG**: пошаговое руководство по созданию ваших первых рабочих процессов - **Отладка и мониторинг**: как понимать статусы задач и быстро находить проблемы -- **Продвинутые возможности**: шаблоны, параметры, управление временем и другие мощные функции +- **Продвинутые возможности**: шаблоны, параметры, управление временем, пулы, XCom, TaskGroup и алертинг ## 📖 Содержание курса @@ -67,10 +67,11 @@ Apache Airflow — это мощный оркестратор рабочих п - Работа с временными зонами - Понимание execution_date и других временных концепций -### [09. Продвинутые возможности Airflow](09%20-%20Продвинутые%20возможности%20Airflow.md) -- Расширенные операторы и сенсоры -- Обработка ошибок и повторные попытки -- Оптимизация и масштабирование DAG +### [09. Пулы, XCom, TaskGroup и алертинг в Airflow](09%20-%20Пулы,%20XCom,%20TaskGroup%20и%20алертинг%20в%20Airflow.md) +- Управление ресурсами с помощью пулов задач +- Обмен данными между задачами через XCom +- Логическая группировка задач с TaskGroup +- Настройка системы оповещений и уведомлений ## 🛠️ Практический стенд