Обзор изменений в Airflow 3

В этой статье описаны основные изменения, представленные в Airflow 3, и их влияние на сервис Airflow в ADO 3.0.0.

ADO предоставляет возможности обновления до Airflow 3, включая поставку пакетов Airflow, миграцию базы данных метаданных, миграцию конфигурации и среды выполнения, проверку обновления и поддержку отката изменений (rollback).

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

Изменения архитектуры

Airflow 3 представляет сервисно-ориентированную архитектуру с более строгой изоляцией между выполнением задач и БД метаданных Airflow.

Web Server заменен на API Server, а компонент DAG Processor стал обязательным. Воркеры взаимодействуют с API Server вместо прямого доступа к БД метаданных.

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

Совместимость DAG

Перед обновлением необходимо проверить совместимость DAG с Airflow 3. Airflow предоставляет правила Ruff для обнаружения несовместимого кода.

Более подробно о том, как обновить DAG для работы с Airflow 3, можно прочитать в статье Migration helper.

Импорты Airflow SDK

Airflow 3 предоставляет airflow.sdk в качестве стабильного интерфейса для DAG. Замените импорты из внутренних модулей Airflow на соответствующие пути airflow.sdk из таблицы ниже.

Airflow 2.x Airflow 3.x

airflow.decorators.dag

airflow.sdk.dag

airflow.decorators.task

airflow.sdk.task

airflow.models.dag.DAG

airflow.sdk.DAG

airflow.models.baseoperator.BaseOperator

airflow.sdk.BaseOperator

airflow.models.param.Param

airflow.sdk.Param

airflow.sensors.base.BaseSensorOperator

airflow.sdk.BaseSensorOperator

airflow.hooks.base.BaseHook

airflow.sdk.BaseHook

airflow.utils.task_group.TaskGroup

airflow.sdk.TaskGroup

airflow.utils.context.Context

airflow.sdk.Context

airflow.datasets.Dataset

airflow.sdk.Asset

airflow.datasets.DatasetAlias

airflow.sdk.AssetAlias

airflow.models.connection.Connection

airflow.sdk.Connection

airflow.models.variable.Variable

airflow.sdk.Variable

Standard provider

Некоторые операторы, сенсоры и триггеры, которые ранее входили в основной пакет Airflow, были перенесены в Standard provider. К ним относятся часто используемые компоненты, такие как BashOperator, PythonOperator, ExternalTaskSensor и FileSensor.

Прямой доступ к БД метаданных

Код задач и пользовательские операторы больше не могут использовать сессии базы данных Airflow для прямого доступа к базе данных метаданных.

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

Прямой доступ к базе данных через хуки БД является только обходным решением для случаев, при которых невозможно обработать данные через API. Такой доступ не рекомендуется, поскольку схема базы данных метаданных не является публичным API и может измениться в будущих версиях Airflow.

Удаленные функции и переменные контекста

В Airflow 3 были убраны некоторые возможности, помеченные как устаревшие в Airflow 2.x. Убедитесь, что эти функции не используются в ваших DAG:

  • SubDAG;

  • SequentialExecutor;

  • CeleryKubernetesExecutor или LocalKubernetesExecutor;

  • SLA;

  • параметр CLI --subdir (-S);

  • REST API /api/v1.

Используйте TaskGroups вместо SubDAG, Deadline Alerts вместо SLA, Multiple Executor Configuration вместо удаленных Kubernetes executor и /api/v2 вместо /api/v1.

Также удалены следующие переменные:

tomorrow_ds
tomorrow_ds_nodash
yesterday_ds
yesterday_ds_nodash
prev_ds
prev_ds_nodash
prev_execution_date
prev_execution_date_success
next_execution_date
next_ds_nodash
next_ds
execution_date

Перед обновлением сервиса отредактируйте DAG, использующие эти переменные.

Изменения планирования DAG

Значение catchup_by_default по умолчанию теперь равно False.

Значение create_cron_data_intervals по умолчанию также равно False. Поэтому DAG, использующие простое cron-выражение, используют CronTriggerTimetable вместо CronDataIntervalTimetable.

Если DAG зависит от значений data_interval_start или data_interval_end, перед обновлением установите для параметра create_cron_data_intervals значение True.

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

Например:

from airflow.sdk import dag, task

@dag
def process_data():
    @task
    def process():
        from airflow.sdk import get_current_context

        context = get_current_context()
        processing_date = context["logical_date"]
        return f"Processing data for {processing_date}"

    process()

process_data()

Если DAG требуется вычисленный интервал данных, можно также продолжать использовать data_interval_start и data_interval_end.

Поведение XCom

Поведение xcom_pull() по умолчанию изменилось. Если task_ids не указан, Airflow 3 выполняет поиск только в текущей задаче. Чтобы получить значение XCom из другой задачи, явно укажите task_ids:

value = ti.xcom_pull(task_ids="upstream_task", key="shared_state")

Изменения аутентификации

В ADO в качестве менеджера аутентификации по умолчанию используется FAB, а провайдер FAB поставляется вместе с бандлом ADO.

Теперь маршруты аутентификации имеют префикс /auth. Например, URL перенаправления OAuth изменяется с https://<your-airflow-url.com>/oauth-authorized/google на https://<your-airflow-url.com>/auth/oauth-authorized/google.

Если используется OAuth, OIDC или LDAP, после обновления проверьте работу аутентификации.

Плагины

Плагины, использующие представления или пункты меню Flask-AppBuilder либо blueprints Flask, требуют дополнительной миграции.

Затронутые плагины можно либо перевести на интерфейсы плагинов Airflow 3, либо использовать FAB provider в качестве уровня совместимости.

Предпочтительные интерфейсы Airflow 3 включают:

  • external_views

  • fastapi_apps

  • fastapi_root_middlewares

Обновление в ADO 3.0.0

ADO 3.0.0 поддерживает обновление существующих кластеров ADO с Airflow 2.11.1 до Airflow 3.2.1_arenadata1 без пересоздания кластера.

Процедура обновления включает:

  • поставку пакетов Airflow 3;

  • миграцию БД метаданных;

  • миграцию конфигурации;

  • миграцию среды выполнения и сервисов;

  • проверку обновления;

  • поддержку отката.

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

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

Нашли ошибку? Выделите текст и нажмите Ctrl+Enter чтобы сообщить о ней