Стилистические правки
This commit is contained in:
@@ -1,6 +1,6 @@
|
|||||||
# Продвинутые возможности Airflow для начинающих специалистов
|
# Пулы, XCom, TaskGroup и алертинг в Airflow
|
||||||
|
|
||||||
В этом материале мы познакомимся с расширенными функциями Airflow, которые помогут вам решать нетривиальные задачи при работе со сложными пайплайнами.
|
В этом материале мы познакомимся с управлением ресурсами через пулы задач, обменом данными между задачами через XCom, группировкой задач с помощью TaskGroup и настройкой системы оповещений (алертинга) в Airflow.
|
||||||
|
|
||||||
Начнем с улучшения нашего базового пайплайна — проведем рефакторинг кода для лучшей читаемости и поддержки.
|
Начнем с улучшения нашего базового пайплайна — проведем рефакторинг кода для лучшей читаемости и поддержки.
|
||||||
|
|
||||||
@@ -12,11 +12,18 @@
|
|||||||
|
|
||||||
Когда планировщик обнаруживает, что наступило время выполнения DAG, он запускает задачу согласно заданной последовательности. При этом задача занимает один слот в пуле и освобождает его после завершения.
|
Когда планировщик обнаруживает, что наступило время выполнения DAG, он запускает задачу согласно заданной последовательности. При этом задача занимает один слот в пуле и освобождает его после завершения.
|
||||||
|
|
||||||
Создать новый пул задач и установить его размер можно через веб-интерфейс Airflow:
|
Создать новый пул задач и установить его размер можно через веб-интерфейс 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"`, и они будут занимать слоты именно этого пула.
|
||||||
|
|
||||||
*Создание нового пула задач в Airflow*
|
|
||||||
|
|
||||||
Зачем нужны пулы задач? Они помогают:
|
Зачем нужны пулы задач? Они помогают:
|
||||||
- Организовать запуск процессов в системе
|
- Организовать запуск процессов в системе
|
||||||
@@ -59,13 +66,34 @@ BashOperator(
|
|||||||
|
|
||||||
💡 **Контекст выполнения** — это набор параметров, передаваемых при запуске задачи, а также метаданные, генерируемые во время выполнения: время старта, время завершения, имя DAG и другие. Контекст можно найти следующим образом:
|
💡 **Контекст выполнения** — это набор параметров, передаваемых при запуске задачи, а также метаданные, генерируемые во время выполнения: время старта, время завершения, имя DAG и другие. Контекст можно найти следующим образом:
|
||||||
|
|
||||||
{: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 }}`.
|
||||||
|
|
||||||

|
В 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 предназначен для обмена небольшими сообщениями. Данные проходят сериализацию/десериализацию при чтении и записи в таблицу.
|
**Важное правило**: XCom предназначен для обмена небольшими сообщениями. Данные проходят сериализацию/десериализацию при чтении и записи в таблицу.
|
||||||
|
|
||||||
@@ -73,7 +101,7 @@ BashOperator(
|
|||||||
|
|
||||||
Для передачи больших объемов данных используйте внешние средства: файловую систему, базы данных (чаще всего PostgreSQL) или очереди сообщений (например, Kafka).
|
Для передачи больших объемов данных используйте внешние средства: файловую систему, базы данных (чаще всего PostgreSQL) или очереди сообщений (например, Kafka).
|
||||||
|
|
||||||
Механизм XCom похож на работу функций в Python. Многие операторы по умолчанию возвращают результат выполнения задачи. За это отвечает параметр `do_xcom_push`, который по умолчанию равен `True`.
|
Механизм XCom похож на работу функций в Python. Многие операторы (например, `PythonOperator`) по умолчанию возвращают результат выполнения задачи. За это отвечает параметр `do_xcom_push`, который во многих случаях равен `True` по умолчанию.
|
||||||
|
|
||||||
Чтобы прочитать сообщения из XCom, используйте метод `xcom_pull` в контексте задачи:
|
Чтобы прочитать сообщения из XCom, используйте метод `xcom_pull` в контексте задачи:
|
||||||
|
|
||||||
@@ -149,17 +177,7 @@ with TaskGroup("data_loading") as loading_group:
|
|||||||
|
|
||||||
Обратите внимание: при использовании TaskGroup последовательность задач указывается внутри группы после объявления всех задач, а в конце DAG описывается последовательность выполнения самих групп.
|
Обратите внимание: при использовании TaskGroup последовательность задач указывается внутри группы после объявления всех задач, а в конце DAG описывается последовательность выполнения самих групп.
|
||||||
|
|
||||||
Визуально в интерфейсе Airflow это выглядит так:
|
Визуально в интерфейсе Airflow группы задач отображаются как один узел с небольшим индикатором. Клик по нему разворачивает или сворачивает вложенные задачи, что значительно улучшает восприятие DAG с большим количеством задач и связей, особенно когда в них десятки и сотни задач.
|
||||||
|
|
||||||

|
|
||||||
|
|
||||||
*Группировка задач с помощью TaskGroup*
|
|
||||||
|
|
||||||
TaskGroup добавляет интерактивность в веб-интерфейс — группы задач можно сворачивать и разворачивать, что значительно улучшает восприятие DAG с большим количеством задач и связей:
|
|
||||||
|
|
||||||

|
|
||||||
|
|
||||||
*Группу задач можно раскрыть для детального просмотра*
|
|
||||||
|
|
||||||
TaskGroup — это удобный способ логической группировки задач, который помогает упростить код и представить сложные пайплайны более компактно.
|
TaskGroup — это удобный способ логической группировки задач, который помогает упростить код и представить сложные пайплайны более компактно.
|
||||||
|
|
||||||
@@ -167,7 +185,7 @@ TaskGroup — это удобный способ логической групп
|
|||||||
|
|
||||||
Алертинг — один из ключевых компонентов системы оркестрации, так как важно своевременно получать уведомления об ошибках для их оперативного анализа и решения.
|
Алертинг — один из ключевых компонентов системы оркестрации, так как важно своевременно получать уведомления об ошибках для их оперативного анализа и решения.
|
||||||
|
|
||||||
По умолчанию в Airflow настроена отправка уведомлений на электронную почту. При создании DAG указываются email-адреса, на которые будут отправляться сообщения. С помощью параметров можно настроить различные сценарии оповещений.
|
В Airflow есть встроенная поддержка отправки уведомлений на электронную почту (при условии, что в конфигурации настроен SMTP-сервер). При создании DAG указываются email-адреса, на которые будут отправляться сообщения. С помощью параметров можно настроить различные сценарии оповещений.
|
||||||
|
|
||||||
Давайте модифицируем наш первый DAG так, чтобы получать уведомления на почту при возникновении ошибок. При этом настроим перезапуск задач в случае неудачи (например, 2 попытки), но без уведомлений о самих перезапусках:
|
Давайте модифицируем наш первый DAG так, чтобы получать уведомления на почту при возникновении ошибок. При этом настроим перезапуск задач в случае неудачи (например, 2 попытки), но без уведомлений о самих перезапусках:
|
||||||
|
|
||||||
@@ -208,14 +226,16 @@ dag = DAG(
|
|||||||
from datetime import datetime, timedelta, timezone
|
from datetime import datetime, timedelta, timezone
|
||||||
import dateutil
|
import dateutil
|
||||||
from airflow.utils.email import send_email_smtp
|
from airflow.utils.email import send_email_smtp
|
||||||
|
from airflow.operators.python import get_current_context
|
||||||
|
|
||||||
MAIL_LIST = [
|
MAIL_LIST = [
|
||||||
"email_1@gmail.ru",
|
"email_1@gmail.ru",
|
||||||
"email_2@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.
|
# 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_start = dateutil.parser.isoparse(dagrun_begin_time.replace('Z', '+00:00'))
|
||||||
calculation_end = datetime.now(timezone.utc)
|
calculation_end = datetime.now(timezone.utc)
|
||||||
duration = str(timedelta(seconds=(calculation_end - calculation_start).seconds))
|
duration = str(timedelta(seconds=(calculation_end - calculation_start).seconds))
|
||||||
@@ -241,6 +261,7 @@ def notify_email(calculation_dt: str, dagrun_begin_time, **context):
|
|||||||
|
|
||||||
```python
|
```python
|
||||||
from airflow.operators.python import PythonOperator
|
from airflow.operators.python import PythonOperator
|
||||||
|
from airflow.utils.trigger_rule import TriggerRule
|
||||||
from utils import notify_email
|
from utils import notify_email
|
||||||
|
|
||||||
LOCAL_CALCULATION_DT = '{{ dag.timezone.convert(execution_date).strftime("%Y-%m-%d") }}'
|
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(
|
email_notification_task = PythonOperator(
|
||||||
task_id="send_email_notification",
|
task_id="send_email_notification",
|
||||||
python_callable=notify_email,
|
python_callable=notify_email,
|
||||||
provide_context=True,
|
|
||||||
dag=dag,
|
dag=dag,
|
||||||
trigger_rule=TriggerRule.ALL_DONE,
|
trigger_rule=TriggerRule.ALL_DONE,
|
||||||
op_args=[LOCAL_CALCULATION_DT, DAG_RUN_BEGIN_TIME],
|
op_args=[LOCAL_CALCULATION_DT, DAG_RUN_BEGIN_TIME],
|
||||||
@@ -267,6 +287,7 @@ email_notification_task = PythonOperator(
|
|||||||
В этом примере используется переменная 'database_connection_string', предварительно созданная в интерфейсе Airflow, для хранения строки подключения к базе данных. Это позволяет избежать жесткого кодирования конфиденциальной информации в коде DAG и упрощает настройку подключения для разных окружений.
|
В этом примере используется переменная 'database_connection_string', предварительно созданная в интерфейсе Airflow, для хранения строки подключения к базе данных. Это позволяет избежать жесткого кодирования конфиденциальной информации в коде DAG и упрощает настройку подключения для разных окружений.
|
||||||
|
|
||||||
```python
|
```python
|
||||||
|
import os
|
||||||
import datetime as dt
|
import datetime as dt
|
||||||
import pandas as pd
|
import pandas as pd
|
||||||
from airflow.models import DAG
|
from airflow.models import DAG
|
||||||
@@ -351,6 +372,7 @@ start_task >> data_processing
|
|||||||
Итоговый DAG будет выглядеть следующим образом:
|
Итоговый DAG будет выглядеть следующим образом:
|
||||||
|
|
||||||
```python
|
```python
|
||||||
|
import os
|
||||||
import datetime as dt
|
import datetime as dt
|
||||||
import pandas as pd
|
import pandas as pd
|
||||||
from airflow.models import DAG
|
from airflow.models import DAG
|
||||||
|
|||||||
Binary file not shown.
|
Before Width: | Height: | Size: 52 KiB |
Binary file not shown.
|
Before Width: | Height: | Size: 1.3 MiB |
Binary file not shown.
|
Before Width: | Height: | Size: 1.4 MiB |
Binary file not shown.
|
Before Width: | Height: | Size: 9.6 MiB |
Binary file not shown.
|
Before Width: | Height: | Size: 89 KiB |
Reference in New Issue
Block a user