从 Airflow 2 升级到 Airflow 3

Apache Airflow 3 是包含破坏性变更的大版本。本指南按步骤介绍如何从 Airflow 2.x 升级到 Airflow 3.0,并说明迁移时需要适配的架构与行为变化。

理解 Airflow 3.x 的架构变化

Airflow 3.x 在安全性、可扩展性和可维护性方面做出了重要架构调整。了解这些变化,有助于在升级前准备并调整工作流。

Airflow 2.x 架构

Airflow 2.x architecture diagram showing scheduler, metadata database, and worker
  • 所有组件直接与 Airflow 元数据库通信。
  • Airflow 2 的设计假设各组件处在同一网络空间。任务代码和执行任务的 Airflow 包代码运行在同一进程中。
  • worker 直接访问 Airflow 数据库,并执行全部用户代码。
  • 用户代码能够导入数据库会话,并可能对元数据库执行恶意操作。
  • 数据库连接数量容易过多,带来扩展困难。

Airflow 3.x 架构

Airflow 3.x architecture diagram showing the decoupled Execution API Server and worker subprocesses
  • 对任务和 worker 而言,API 服务器现在是访问元数据库的唯一入口。
  • API 服务器承载多个应用:Airflow REST API、为界面提供支持并托管静态 JavaScript 的内部 API,以及 worker 执行任务实例时使用的任务执行 API。
  • worker 与 API 服务器通信,而不再直接连接数据库。
  • DAG Processor 与 Triggerer 在执行任务、特别是需要变量或连接时,也使用任务执行机制。

数据库访问限制

Airflow 3 限制任务代码直接访问元数据库,这是重要的安全与架构改进,也改变了 DAG 作者访问 Airflow 资源的方式。

  • 不再允许任务代码直接导入并使用 Airflow 数据库会话或模型。
  • 状态转换、心跳、XCom 和资源获取等运行时交互,通过专门的 Task Execution API 完成。
  • worker 内的任务代码不能直接访问或修改元数据库,从而改善隔离与安全。但 DAG 作者的代码在 DAG File Processor 和 Triggerer 中仍可能以能够直接访问数据库的方式运行,详见 Airflow 安全模型。
  • Task SDK 为访问 Airflow 资源提供稳定、向前兼容的接口,减少对数据库结构的直接依赖。

第一步:满足前置条件

  • 当前 Airflow 至少应为 2.7,推荐先升级到最新的 2.x,再升级到 Airflow 3。
  • 确保 Python 版本在支持范围内。
  • 确保没有继续使用 Airflow 3 已移除的功能,相关清单见后文。

第二步:清理并备份现有实例

开始迁移前,强烈建议备份 Airflow 实例,尤其是元数据库。

如果数据库不支持热备份,应先关闭 Airflow 实例,再备份,以保证数据一致性。否则备份可能无法包含全部 TaskInstance 或 DagRun。

如果未备份而迁移失败,例如 Airflow CLI 与数据库之间的网络连接中断,数据库可能停留在部分迁移状态。备份是应对此类故障的重要保障。

长期运行的实例可能积累大量不再需要的数据,例如旧 XCom。Airflow 3 升级包含数据库结构变更,数据库越大,迁移可能越慢。建议事先清理元数据库,可使用 airflow db clean 缩减数据规模,以提升迁移速度和安全性。

还应消除 DAG 处理错误,例如 AirflowDagDuplicatedIdException。airflow dags reserialize 应能无错误运行。如果需要修改 DAG,先将修改部署到旧实例,等待所有 DAG 重新处理且错误全部消失,再进行升级。

第三步:DAG 作者检查兼容性

Airflow 项目基于 Ruff 及其 AIR 规则提供了升级检查能力。

AIR301 和 AIR302 指示 Airflow 3 的破坏性变更;AIR311 和 AIR312 则指出当前尚未造成破坏、但强烈建议更新的用法。最新 Ruff 包含最新规则,至少应使用 0.13.1。

检查需要在 Airflow 3 上修复的 DAG 不兼容项:

ruff check dags/ --select AIR301

预览建议修改:

ruff check dags/ --select AIR301 --show-fixes

自动修复能够安全处理的变化:

ruff check dags/ --select AIR301 --fix

部分修复标记为 unsafe。它们通常不会破坏 DAG 代码,但可能改变运行时行为,因此被标记为不安全。说明见 Fix Safety。启用这类修复:

ruff check dags/ --select AIR301 --fix --unsafe-fixes

在 AIR 规则中,保持导入成员名不变、只改变路径的修复会被视为 unsafe。例如,将 from airflow.sensors.base_sensor_operator import BaseSensorOperator 改为 from airflow.sdk.bases.sensor import BaseSensorOperator,需要先删除旧导入,再加入新导入。

相比之下,同时改变成员名和路径的修复可以是 safe,例如把 from airflow.datasets import Dataset 改为 from airflow.sdk import Asset,因为不必先删除旧导入。要清理未使用的旧导入,还需启用 unused-import(F401) 规则。

这些选项也可以写入配置文件,详见 Ruff 配置。

主要导入路径变化

Ruff 能自动处理许多导入问题,但应了解以下迁移对应关系。旧路径已弃用,将在未来 Airflow 版本中移除。

旧导入路径(已弃用) 新导入路径
airflow.decorators.dag airflow.sdk.dag
airflow.decorators.task airflow.sdk.task
airflow.decorators.task_group airflow.sdk.task_group
airflow.decorators.setup airflow.sdk.setup
airflow.decorators.teardown airflow.sdk.teardown
airflow.models.dag.DAG airflow.sdk.DAG
airflow.models.baseoperator.BaseOperator airflow.sdk.BaseOperator
airflow.models.param.Param airflow.sdk.Param
airflow.models.param.ParamsDict airflow.sdk.ParamsDict
airflow.models.baseoperatorlink.BaseOperatorLink airflow.sdk.BaseOperatorLink
airflow.sensors.base.BaseSensorOperator airflow.sdk.BaseSensorOperator
airflow.hooks.base.BaseHook airflow.sdk.BaseHook
airflow.notifications.basenotifier.BaseNotifier airflow.sdk.BaseNotifier
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.datasets.DatasetAll airflow.sdk.AssetAll
airflow.datasets.DatasetAny airflow.sdk.AssetAny
airflow.models.connection.Connection airflow.sdk.Connection
airflow.models.variable.Variable airflow.sdk.Variable
airflow.io.* airflow.sdk.io.*

原文给出的迁移时间线是:Airflow 3.1 中旧导入仍可工作,但会显示弃用警告;未来版本将移除这些导入。

第四步:安装 Standard Provider

一些原本随 airflow-core 提供的常用 Operator、Sensor 和 Trigger,已经拆分到独立的 apache-airflow-providers-standard 包中,例如 BashOperator、PythonOperator、ExternalTaskSensor 和 FileSensor。

该包也可安装到 Airflow 2.x,便于提前修改 DAG,让它们从 Standard Provider 导入这些组件,而不是从 Airflow Core 导入。

第五步:检查自定义任务是否直接访问数据库

Airflow 3 中,Operator 不能再通过数据库会话直接访问元数据库。自定义 Operator 应检查并移除这类调用,可参考 issue 49187 中的示例。

以前直接访问元数据库的 Operator 或任务代码,需要迁移到以下方式之一。

推荐方式:Airflow Python Client

使用官方 Airflow Python Client,通过 REST API 与 Airflow 资源交互。它为多数场景提供了 API,包括 DagRun、TaskInstance、Variable、Connection 和 XCom。

优点包括:

  • worker 无需直接访问数据库网络。
  • 符合 Airflow 3 的 API 优先架构。
  • worker 环境不需要数据库凭据,而是使用 API token。
  • worker 无需安装数据库驱动。
  • 访问控制与认证集中在 API 服务器。

需要考虑的限制包括:

  • 需要安装 apache-airflow-client。
  • 需要调用 /auth/token 获取访问令牌,并按需轮换。
  • API 服务器必须可用,而且 worker 必须能通过网络访问它。
  • 并非所有数据库操作都有对应 API。

如果 Python Client 尚不支持所需功能,应考虑请求增加 API 端点或 Task SDK 能力。Airflow 社区优先补足 API,而不是恢复任务直接访问数据库的方式。

已知替代方案:DbApiHook

警告:使用 PostgresHook 或 MySqlHook 直接连接元数据库不受推荐。原文仅将其记录为无法使用 Python Client 时的已知临时方案。它限制明显,并会在未来版本失效。

重要风险包括:

  • 原文明确指出该方案会在 Airflow 3.2 及之后的版本中失效,数据库结构变化时需要自行修改代码。
  • 元数据库结构不是公开 API,可以随时变化,不保证提前通知;结构变化可能直接破坏查询。
  • 这种方式破坏任务隔离,与 Airflow 3 的核心设计相违背。
  • 每个任务重新建立独立数据库连接,会恢复类似 Airflow 2 的连接模式,显著改变性能和扩展特性。

只有在 Python Client 无法满足需求、并且充分理解风险时,才可考虑创建一个指向元数据库的 PostgreSQL 或 MySQL 连接,再通过数据库 Hook 查询。

这些 Hook 通过 psycopg2 或 mysqlclient 等驱动直接连接数据库,不经过 API 服务器。

PostgresHook 示例,MySqlHook 的接口也类似:

from airflow.sdk import task
from airflow.providers.postgres.hooks.postgres import PostgresHook


@task
def get_connections_from_db():
    hook = PostgresHook(postgres_conn_id="metadata_postgres")
    records = hook.get_records(sql="""
        SELECT conn_id, conn_type, host, schema, login
        FROM connection
        WHERE conn_type = 'postgres'
        LIMIT 10;
        """)

    return records

如果更倾向于使用 Operator,也可以使用 SQLExecuteQueryOperator:

from airflow.providers.common.sql.operators.sql import SQLExecuteQueryOperator

query_task = SQLExecuteQueryOperator(
    task_id="query_metadata",
    conn_id="metadata_postgres",
    sql="SELECT conn_id, conn_type FROM connection WHERE conn_type = 'postgres'",
    do_xcom_push=True,
)

注意:元数据库连接始终应使用只读凭据,原文还建议采用临时凭据。

第六步:部署负责人升级实例

Airflow 提供配置升级工具,以简化迁移。先运行配置检查:

airflow config update

也可以让工具自动更新配置,使其兼容 Airflow 3:

airflow config update --fix

数据库升级是整个过程的主要部分。Airflow 3 的数据库迁移流程与 Airflow 2.7 及之后版本一致:

airflow db migrate

如果插件使用 Flask-AppBuilder 视图 appbuilder_views、菜单项 appbuilder_menu_items 或 Flask 蓝图 flask_blueprints,需要将其改为 FastAPI 应用,或者安装提供 Airflow 3 向后兼容层的 FAB provider。

理想做法是迁移到 Airflow 3 插件接口,即外部视图 external_views、FastAPI 应用 fastapi_apps 和中间件 fastapi_root_middlewares。

通过 Airflow Helm Chart 部署时,应将现有 values 与 Airflow 3 可用选项逐项对应。webserver 下的配置需要迁到 apiServer,许多参数已经改名或移除。Chart 专用升级清单还涵盖 values.yaml、独立 DAG Processor、JWT secret、FAB 默认行为、最低 Kubernetes 版本,以及 Chart 1.16.0 到 1.18.0 的键名变化,详见将 Helm Chart 升级到 Airflow 3。

第七步:修改启动脚本

Airflow 3 的 Webserver 已变为通用 API 服务器,启动命令为:

airflow api-server

DAG Processor 现在必须独立启动,即使本地开发环境也一样:

airflow dag-processor

完成后即可启动 Airflow 3 实例。

第八步:升级后需要关注的事项

如果通过 OAuth、OIDC 或 LDAP 配置了单点登录,应确认认证仍按预期工作。

如果使用自定义 webserver_config.py,需要将 from airflow.www.security import AirflowSecurityManager 替换为 from airflow.providers.fab.auth_manager.security_manager.override import FabAirflowSecurityManagerOverride。

破坏性变更

Airflow 2.x 中已弃用的一些能力,在 Airflow 3 中不再可用:

  • SubDAG:由 TaskGroup、Asset 和数据感知调度替代。
  • Sequential Executor:由 LocalExecutor 替代,后者可与 SQLite 配合用于本地开发。
  • CeleryKubernetesExecutor 与 LocalKubernetesExecutor:由多执行器配置替代。
  • SLA:已弃用并移除,由 Deadline Alerts 替代。
  • Subdir:许多 CLI 命令中的 --subdir 或 -S 已由 DAG bundles 替代。
  • REST API /api/v1:由基于 FastAPI 的稳定版 /api/v2 替代,详见 API v2。

以下任务实例上下文变量不再提供,未替换会导致 DAG 错误: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。

catchup_by_default 现在默认为 False。

create_cron_data_intervals 现在也默认为 False,因此默认使用 CronTriggerTimetable,而非 CronDataIntervalTimetable。这只影响直接给 schedule= 传入 cron 字符串的 DAG,例如 schedule="0 0 * * *";显式传入 timetable 实例的 DAG 不受影响。

应确认任务是否依赖 data_interval_start、data_interval_end,以及由 logical_date 派生、在两种 timetable 下可能发生偏移的 ds、ts 等模板值。如果依赖这些区间语义,应明确设置 create_cron_data_intervals=True,保留 CronDataIntervalTimetable;否则新的默认值通常合适。

应在升级前完成这一设置。如果已有 Airflow 3 DagRun 后才从 CronTriggerTimetable 切回 CronDataIntervalTimetable,为避免与上一运行的 logical_date 冲突,会跳过一次计划运行。

手动运行与数据区间

在 Airflow 3 中,不能假设手动触发的 DAG 运行,其 data_interval 一定由传入的 logical_date 推导,或与之相等。如果逻辑需要用户指定的触发日期,应显式使用 logical_date。手动运行或使用 TriggerDagRunOperator 时读取数据区间的工作流,尤其需要注意。

默认认证管理器

默认 auth_manager 现在是 Simple Auth。如果需要继续使用 FAB,应安装 FAB provider,并把认证管理器设置为 FabAuthManager:

airflow.providers.fab.auth_manager.fab_auth_manager.FabAuthManager

认证管理器定义的 API 路由现在统一带 /auth 前缀。应用之外使用的 URL,例如 OAuth 回调地址,也必须相应更新。例如 Airflow 2.x 的 https://<your-airflow-domain>/oauth-authorized/google,在 Airflow 3.x 中变为 https://<your-airflow-domain>/auth/oauth-authorized/google。

XCom pull 的默认行为

调用 xcom_pull() 时不传 task_ids,现在只会读取当前任务的值。Airflow 2 中省略该参数,会在同一 DAG 运行的所有任务中搜索,并返回给定键最近推送的值。要读取其他任务的 XCom,现在必须明确指定 task_ids:

# Airflow 2 - pulls most recent value from any task
value = ti.xcom_pull(key="shared_state")

# Airflow 3 - same call only checks the current task
value = ti.xcom_pull(key="shared_state")

# Airflow 3 - specify task_ids to pull from other tasks
value = ti.xcom_pull(task_ids="upstream_task", key="shared_state")

手动 DAG 运行与 logical_date

对于计划运行,logical_date 和 data_interval 都由 DAG 的 timetable 推导。

对于 Airflow 3 的手动运行,不应假设 data_interval_start 或 data_interval_end 来自传入的 logical_date,或与它相等。最终数据区间取决于 timetable 与触发路径,有些 API 还允许显式指定区间。

以下 DAG 尤其需要检查这一差异:

  • 手动运行时读取 data_interval_start 或 data_interval_end。
  • 通过 TriggerDagRunOperator 触发下游 DAG。
  • 从 Airflow 2 迁移而来,并曾把 data_interval_start 当作用户请求的手动运行日期。

迁移建议

如果业务逻辑需要用户为手动运行指定的日期,应明确使用 logical_date:

from airflow.decorators import get_current_context, task


@task
def process_data():
    context = get_current_context()
    processing_date = context["logical_date"]
    return f"Processing data for {processing_date}"

如果业务逻辑需要的是运行最终解析出的区间语义,则继续使用 data_interval_start 和 data_interval_end。

从 Airflow 2 升级时,应梳理手动触发且读取这两个区间字段的工作流,明确它们真正需要的是数据区间,还是用户请求的逻辑日期。

原文来源:Upgrading to Airflow 3。本文依据留存原文译为中文,代码示例按原文保留。

© The Apache Software Foundation. Apache Airflow、Apache、Airflow 及相关标志是 Apache Software Foundation 的注册商标或商标,其他产品及品牌归各自所有者所有。原文许可:License;Apache 许可证。

© 版权声明
THE END
喜欢就支持一下吧
点赞0 分享
评论 抢沙发

请登录后发表评论

    暂无评论内容