Files
airflow-manual/07 - Шаблоны, переменные и подключения в Airflow.md

14 KiB
Raw Permalink Blame History

Шаблоны, переменные и подключения в Airflow

В этом материале вы познакомитесь с мощными инструментами Airflow для создания гибких и переиспользуемых пайплайнов: динамическими шаблонами, безопасными переменными и централизованными подключениями к внешним системам.

Динамические шаблоны Airflow (на основе Jinja)

В процессах обработки данных часто возникает необходимость использовать переменные значения, такие как дата выполнения задачи, для фильтрации исходных данных. В сложных сценариях может потребоваться информация о предыдущих успешных запусках или метаданные текущего выполнения. Airflow предоставляет встроенный механизм динамических шаблонов, основанный на языке Jinja, который значительно упрощает работу с такими переменными.

Синтаксис шаблонов интуитивно понятен — переменная заключается в двойные фигурные скобки с пробелами по краям. Для эффективного использования достаточно знать, какое значение возвращает конкретный шаблон, и разместить его в нужном месте кода.

Подробный справочник доступных шаблонов можно найти в официальной документации Airflow.

Вот практические примеры использования:

Передача даты выполнения в задачу:

В этом примере создается DAG с идентификатором "template_example", который запускается каждые 15 минут. Задача "display_date" использует шаблон {{ ds }} для вывода текущей даты выполнения в формате YYYY-MM-DD.

...

dag = DAG(
    dag_id="template_example",
    schedule_interval="*/15 * * * *",
    default_args=default_args
) 

t1 = BashOperator(task_id="display_date", bash_command="echo {{ ds }}")

t1

Использование шаблонов для лучшей отслеживаемости задач:

В этом примере создается DAG с идентификатором "template_tracking_example", который запускается каждые 20 минут. Вторая задача использует шаблон {{ ds }} в команде bash, чтобы явно указывать дату обработки в логах.

...

dag = DAG(
    dag_id="template_tracking_example",
    schedule_interval="*/20 * * * *",
    default_args=default_args
) 

t1 = BashOperator(task_id="show_date", bash_command="echo {{ ds }}")
t2 = BashOperator(task_id="process_for_date", bash_command="echo Processing for {{ ds }}")

t1 >> t2

Безопасные переменные (Variables)

Переменные Airflow представляют собой пары "ключ-значение", хранящиеся в метадатабазе системы. Они идеально подходят для хранения конфигурационных параметров, таких как пути к скриптам, имена таблиц или другие настройки, которые должны быть доступны в разных DAG.

Управление переменными осуществляется через веб-интерфейс Airflow (раздел Admin → Variables). Через этот раздел можно:

  • создавать и редактировать пары «ключ-значение» вручную;
  • импортировать набор переменных из JSON-файла;
  • удалять больше не нужные настройки.

Как создать переменную через UI

Интерфейс ниже соответствует Airflow 2.9.x:

  1. Откройте веб-интерфейс Airflow и авторизуйтесь под пользователем с правами Admin.

  2. В верхнем меню выберите Admin → Variables.

  3. В правом верхнем углу нажмите кнопку + Add a new record (или иконку +).

  4. В поле Key задайте имя переменной.
    Например, создадим переменную с паролем к учебной БД отчётности PostgreSQL:

    • Key: reporting_db_password
  5. В поле Value введите значение.
    Например:

    • Value: airflow_report_ro
  6. Поле Description можно использовать для короткого пояснения, зачем нужна переменная, например:
    Пароль read-only к учебной БД отчётности.

  7. Нажмите Save.

После сохранения переменная появится в таблице. Значение будет частично скрыто в UI: вместо реального пароля вы увидите *** — Airflow маскирует секреты в интерфейсе и логах, чтобы их нельзя было случайно подсмотреть.

Теперь эту переменную можно использовать в коде DAG, например:

from airflow.models import Variable

reporting_db_password = Variable.get("reporting_db_password")

Для защиты конфиденциальной информации Airflow автоматически маскирует значения переменных, в названии которых содержится слово secret, а также ряд других чувствительных паттернов.

Подробнее о переменных — в официальной документации Airflow.

Пример использования переменной в коде DAG:

В этом примере создается DAG с идентификатором "variable_example", который использует переменную 'data_storage_path', предварительно сохраненную в Airflow. Значение переменной извлекается с помощью Variable.get() и используется в команде bash для указания пути к данным.

from airflow import DAG
from airflow.operators.bash import BashOperator
from airflow.models import Variable
from datetime import datetime

default_args = {
    'owner': 'data_team',
    'start_date': datetime(2023, 1, 1),
    'retries': 1,
}

dag = DAG(
    dag_id="variable_example",
    schedule_interval=None,
    default_args=default_args
)

# Получение значения переменной
data_path = Variable.get('data_storage_path')

task = BashOperator(
    task_id='process_data',
    bash_command=f'echo "Processing data from {data_path}"',
    dag=dag
)

Переменные делают код DAG более читаемым и модульным, позволяя легко адаптировать один и тот же пайплайн для разных окружений или сценариев использования.

Централизованные подключения (Connections)

Подключения в Airflow — это безопасный способ хранения учетных данных и параметров для взаимодействия с внешними системами. Каждое подключение имеет уникальный идентификатор (conn_id) и содержит необходимые параметры: хост, порт, логин, пароль и другие специфичные настройки.

Подключения поддерживают широкий спектр систем:

  • Базы данных (PostgreSQL, MySQL, Oracle и др.)
  • Облачные хранилища (AWS S3, Google Cloud Storage)
  • Системы уведомлений (Email, Telegram)
  • И многие другие через Airflow Providers

При использовании операторов, взаимодействующих с внешними системами, Airflow автоматически ищет соответствующее подключение с суффиксом _default. Однако можно явно указать альтернативное подключение:

В приведенном примере используется PostgresOperator для создания таблицы в базе данных. Вместо использования подключения по умолчанию, явно указывается подключение с идентификатором 'my_postgres_conn'.

from airflow.providers.postgres.operators.postgres import PostgresOperator

create_table = PostgresOperator(
    task_id='create_user_table',
    sql='''
        CREATE TABLE users(
        user_id integer NOT NULL,
        created_at TIMESTAMP NOT NULL
        );''',
    postgres_conn_id='my_postgres_conn'
)

Управление подключениями доступно через интерфейс Airflow (Admin → Connections). Если требуемый тип подключения отсутствует, его можно добавить установкой соответствующего Airflow Provider из официального репозитория.

Управление подключениями доступно через интерфейс Airflow (раздел Admin → Connections). Подключения хранятся в метадатабазе Airflow, а пароли и другие чувствительные поля шифруются с помощью Fernet и маскируются в UI и логах.

Как создать подключение к PostgreSQL через UI

Интерфейс ниже соответствует Airflow 2.9.x и стандартному Docker-стенду из документации:

  1. Откройте веб-интерфейс Airflow.

  2. В верхнем меню выберите Admin → Connections.

  3. В правом верхнем углу нажмите кнопку + Add a new record.

  4. В форме укажите параметры:

    • Connection Id: my_postgres_conn
      Это имя мы будем использовать в коде DAG (параметр postgres_conn_id).
    • Connection Type: Postgres
    • Host: postgres
      (так называется контейнер PostgreSQL в типовом docker-compose.yaml из официальной инструкции).
    • Schema: airflow
      (имя базы данных; в учебном стенде можно использовать стандартную БД).
    • Login: airflow
    • Password: airflow
    • Port: 5432
  5. Нажмите кнопку Test (если доступна) — Airflow попробует подключиться к базе.

  6. Если тест успешен, нажмите Save.

Теперь подключение с идентификатором my_postgres_conn доступно во всех DAG’ах. В примере ниже PostgresOperator явно использует это подключение:

from airflow.providers.postgres.operators.postgres import PostgresOperator

create_table = PostgresOperator(
    task_id="create_user_table",
    sql="""
        CREATE TABLE users(
            user_id    INTEGER      NOT NULL,
            created_at TIMESTAMP    NOT NULL
        );
    """,
    postgres_conn_id="my_postgres_conn",
)

Для успешной работы с внешними системами сначала необходимо создать соответствующее подключение, а затем использовать его идентификатор в операторах вашего DAG.

Более подробную информацию о настройке подключений можно найти в документации Airflow.

Проверочный список для качественного DAG

После создания DAG задайте себе следующие вопросы для обеспечения его качества и безопасности:

  • Сможет ли коллега понять и поддерживать этот DAG в моё отсутствие?
  • Содержит ли код чувствительную информацию (логины, пароли, API-ключи)?
  • Какие параметры можно вынести в переменные для лучшей гибкости?
  • Требуется ли маскировка конфиденциальных значений?
  • Используются ли в логике даты или временные метки, которые можно заменить на шаблоны?

Ответы на эти вопросы помогут вам создавать надежные, безопасные и легко поддерживаемые пайплайны в Airflow.