Apache Airflow 3 升级全指南:从 2.x 迁移的架构变化、破坏性变更与分步实操
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
Apache Airflow 3 是 Airflow 项目的一个主版本(major release),包含大量破坏性变更(breaking changes)。本文以仓库中的官方升级指南 upgrading_to_airflow3.rst 为主体,结合airflow-core、task-sdk等子仓库的源码实现,系统讲解从 Airflow 2.x 升级到 3.0 的完整路径:先理解架构变化的底层原因,再按 8 个步骤完成备份、Dag 兼容性检查、配置迁移与数据库升级,最后梳理需要重点排查的破坏性变更清单。读完本文,你将能够独立规划并执行一次安全、可控的 Airflow 3 升级。
理解 Airflow 3.x 的架构变化
Airflow 3.x 引入了显著的安全、可扩展性与可维护性改进。理解这些变化是升级准备的第一步,它决定了后续每一步操作(尤其是 Dag 兼容性改造)的方向。
Airflow 2.x 架构:全员直连数据库
在 Airflow 2.x 中:
- 所有组件(Scheduler、Worker、Triggerer、Webserver 等)都直接与 Airflow 元数据库(metadata database)通信;
- Airflow 2 被设计为让所有组件运行在同一个网络空间内:任务代码与执行任务的 Airflow 包代码运行在同一个进程里;
- Worker 直接连接 Airflow 数据库并执行所有用户代码;
- 由于用户代码可以直接 import 数据库会话(session),存在对元数据库执行恶意操作的风险;
- 各组件对数据库的连接数量过大,带来了明显的扩展性挑战。
简言之,2.x 的信任边界是"组件与用户代码都在数据库旁边",安全问题与连接数问题都由此而来。
Airflow 3.x 架构:API Server 成为唯一入口
Airflow 3.x 的核心变化是引入了解耦的执行 API Server:
- API Server 目前是任务(tasks)和 Worker 访问元数据库的唯一入口,承载了多种应用形态:Airflow REST API、为 Airflow UI(托管静态 JS)服务的内部 API、以及 Worker 通过任务执行接口(task execution interface)执行 TaskInstance(TI)时交互所用的 API;
- Worker 不再直连数据库,而是与 API Server 通信;
- Dag Processor 与 Triggerer 也通过任务执行机制运行它们的任务,尤其是在需要变量(variables)或连接(connections)时。
这一设计把"谁能访问数据库"收敛为"只有 API Server 能访问数据库",从架构层面隔离了用户代码与元数据库。
数据库访问限制:Dag 作者必须知晓的安全边界
Airflow 3 对任务代码直接访问元数据库做了严格限制,这是最影响 Dag 作者的关键变化:
- 禁止直接访问数据库:任务代码不能再直接 import 并使用 Airflow 的数据库会话(sessions)或模型(models);
- 基于 API 的资源访问:所有运行时交互(状态转换、心跳 heartbeat、XCom 以及资源获取)都通过专用的 Task Execution API 完成;
- 安全性增强:通过阻止 Worker 任务代码直接访问或修改元数据库来提升隔离性与安全性。注意,Dag 作者代码在 Dag File Processor 和 Triggerer 中仍可能以直接数据库访问方式执行,详见 security_model.rst;
- 稳定的接口:Task SDK 提供了稳定且向前兼容的接口来访问 Airflow 资源,无需依赖数据库。
从源码上看,Task SDK 正是这一架构落地的载体。在 task-sdk/src/airflow/sdk/init.py 中可以看到,airflow.sdk包通过__lazy_imports惰性导出了DAG、BaseOperator、BaseHook、BaseSensorOperator、BaseNotifier、Connection、Variable、Param、ParamsDict、TaskGroup、Context、Asset、AssetAlias、AssetAll、AssetAny、dag/task/task_group/setup/teardown装饰器以及get_current_context等全部核心符号——这些正是升级后 Dag 代码的新导入目标(详见下文"关键导入路径更新")。
Step 1:检查升级前置条件
开始迁移之前,请确认:
- 当前版本必须是 Airflow 2.7 或更高版本。官方推荐先升级到最新的 2.x 版本,再升级到 Airflow 3;
- Python 版本必须在受支持列表中,先确认你的运行环境满足 Airflow 3 的 Python 版本要求;
- 确认没有使用任何已在 Airflow 3 中移除的功能(完整清单见下文"破坏性变更"一节)。
Step 2:清理并备份现有 Airflow 实例
数据库迁移是升级中风险最高的一环,备份与清理是必须的前置动作:
- 强烈建议在迁移前备份 Airflow 实例,尤其是元数据库:
- 如果你的数据库不具备"热备份(hot backup)"能力,应在关闭 Airflow 实例之后再备份,以保证备份的一致性。否则(例如不关闭实例),备份将不包含所有 TaskInstance 或 DagRun;
- 如果没有备份而迁移失败,可能会进入"半迁移"状态——例如迁移过程中 Airflow CLI 与数据库之间的网络连接中断就可能造成这种情况。备份是避免此类问题的关键预防措施;
- 清理元数据库:长期运行的实例会积累大量不再需要的数据(例如旧 XCom 数据)。Airflow 3 升级过程包含 schema 变更,数据库越大迁移耗时越长。为了更快、更安全地迁移,建议在升级前用
airflow db clean命令(对应 CLI 定义位于 cli_config.py 中的db clean子命令)清理 Airflow 数据库; - 确认 Dag 处理无错误:确保不存在诸如
AirflowDagDuplicatedIdException之类的 Dag 处理错误,应能够无错误地运行airflow dags reserialize。如果存在需要解决的 Dag 处理错误,请先在你的旧实例上部署修复,并等待所有 Dag 重新处理完毕、错误全部消失后,再进行升级。
Step 3:Dag 作者——检查 Dag 的兼容性
为了最小化升级摩擦,Airflow 社区基于Ruff与AIR(Airflow)规则创建了 Dag 升级检查工具。规则 AIR301 与 AIR302 标记的是 Airflow 3 中的破坏性变更;AIR311 与 AIR312 标记的则是当前尚未破坏、但强烈建议更新的变更。
请使用最新的ruff版本(至少为 0.13.1)来获得最新规则。以下命令用于检查 dags 目录中需要在 Airflow 3 上修复才能正常工作的不兼容问题:
ruff check dags/ --select AIR301预览推荐的修复:
ruff check dags/ --select AIR301 --show-fixes部分变更可以自动修复:
ruff check dags/ --select AIR301 --fix部分修复被标记为unsafe(不安全)。不安全修复通常不会破坏 Dag 代码,标记为不安全是因为它们可能改变某些运行时行为。要触发这类修复,使用:
ruff check dags/ --select AIR301 --fix --unsafe-fixes关于 AIR 规则中安全/不安全修复的区别:不安全修复涉及在保持导入成员名称不变的情况下修改导入路径,例如把from airflow.sensors.base_sensor_operator import BaseSensorOperator改为from airflow.sdk.bases.sensor import BaseSensorOperator,这要求 ruff 先删除原导入再添加新导入;而安全修复则是同时修改成员名称与导入路径,例如把from airflow.datasets import Dataset改为from airflow.sdk import Asset,这类调整不需要 ruff 删除旧导入。要清理未使用的遗留导入,需要启用unused-import规则(F401)。这些标记同样可以通过 Ruff 配置文件进行配置。
关键导入路径更新
虽然 ruff 可以自动修复大量导入问题,但下面这张对照表是 Dag 及其他代码在 Airflow 3 中正确导入组件时需要手工掌握的核心变更。旧路径已被弃用,将在未来的 Airflow 版本中移除:
| 旧导入路径(已弃用) | 新导入路径(airflow.sdk) |
|---|---|
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:遗留导入会显示弃用警告,但继续可用;
- 未来的 Airflow 版本:遗留导入将被彻底移除。
这些新路径与 task-sdk/src/airflow/sdk/init.py 中导出的符号一一对应,读者可以在该文件中核对每个符号的实际归属模块。
Step 4:安装 Standard Provider
- 一些原本随
airflow-core包捆绑的常用 Operators、Sensors 与 Triggers(例如BashOperator、PythonOperator、ExternalTaskSensor、FileSensor等)已被拆分到独立的apache-airflow-providers-standard包; - 方便的是,这个包也可以安装在 Airflow 2.x 上,这样 Dag 可以先改为从 standard provider 包引用这些 Operator,而不是从 Airflow Core 引用,从而平滑过渡。
Step 5:审查自定义任务中的直接数据库访问
在 Airflow 3 中,Operator 不能再使用数据库会话直接访问元数据库。如果你有自定义 Operator,请审查代码,确保没有直接数据库访问调用。社区提供了大量修改示例(参见 Apache Airflow issue 49187 中的讨论)。
如果你有自定义 Operator 或任务代码此前直接访问过元数据库,必须迁移到以下方案之一:
推荐方案:使用 Airflow Python Client
使用官方的 Airflow Python Client 通过 REST API 与元数据库交互。Python Client 为大多数用例定义了 API,包括 DagRuns、TaskInstances、Variables、Connections、XComs 等。
优点:
- Worker 无需直接数据库网络访问;
- 与 Airflow 3 的 API-first 架构最契合;
- Worker 环境无需数据库凭据(改用 API 令牌);
- Worker 无需安装数据库驱动;
- 通过 API Server 实现集中式访问控制与认证。
缺点:
- 需要安装
apache-airflow-client包; - 需要通过调用
/auth/token获取访问令牌并按要求轮换; - 依赖 API Server 可用性与到 API Server 的网络连通性;
- 并非所有数据库操作都能通过 API 端点暴露。
注意:如果需要 Python Client 未提供的能力,可以考虑请求新的 API 端点或 Task SDK 功能。Airflow 社区优先补充缺失的 API 能力,而非开放直接数据库访问。
已知的变通方案:使用 DbApiHook(PostgresHook 或 MySqlHook)
警告:此方案不被推荐,仅作为无法使用 Python Client 的用户的已知变通方案被记录。该方案存在显著限制,并且在未来的 Airflow 版本中会失效。
需要重点考虑的事项:
- 未来版本会失效:该方案在 Airflow 3.2+ 及之后会失效,schema 变更时你需要自行负责适配代码;
- 数据库 schema 不是公共 API:Airflow 元数据库 schema 可能随时变更且不另行通知,schema 变更会毫无预兆地破坏你的查询;
- 破坏任务隔离:这与 Airflow 3 的核心特性——任务隔离——相矛盾。任务不应直接访问元数据库;
- 性能影响:这会重新引入 Airflow 2 的行为——每个任务各自打开数据库连接,从根本上改变性能特征与扩展性。
如果你的用例无法通过 Python Client 解决,并且你理解上述风险,可以使用数据库钩子(database hooks)直接查询元数据库。创建一个指向元数据库的数据库连接(PostgreSQL 或 MySQL,与你的元数据库类型一致),然后在 Airflow 中使用 Database Hooks。
注意:这些钩子直接连接数据库(不经由 API Server),使用 psycopg2 或 mysqlclient 等数据库驱动。
使用 PostgresHook 的示例(MySql 也有类似接口):
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使用 SQLExecuteQueryOperator 的示例:
如果你更倾向于使用 Operator 而非 Hook,也可以使用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, )注意:为元数据库连接始终使用只读数据库凭据,并建议使用临时凭据。
Step 6:部署管理人员——升级 Airflow 实例
为让升级过程更简单、更安全,Airflow 3 提供了配置升级工具。先从配置检查开始:
airflow config update该工具也能把配置自动更新为兼容 Airflow 3:
airflow config update --fix从源码看,该命令的实现位于 config_command.py:lint_config函数扫描airflow.cfg中在 Airflow 3.0 里被移除或重命名的参数并给出建议;update_config函数默认执行 dry-run(仅展示变更),--fix才会真正改写airflow.cfg(改写前会先备份旧文件,且会清理原有注释),--all-recommendations则把非破坏性的推荐变更一并纳入。命令内置的CONFIGS_CHANGES列表覆盖了数百项迁移规则,例如:
core.executor默认值由SequentialExecutor改为LocalExecutor;core.sql_alchemy_conn等一批参数从core段迁移到database段;scheduler.catchup_by_default默认值由True改为False;scheduler.create_cron_data_intervals与create_delta_data_intervals默认值均由True改为False;webserver段的web_server_host/web_server_port等参数迁移到api段的host/port;scheduler.max_threads更名为dag_processor.parsing_processes等。
升级的最大组成部分是数据库升级。Airflow 3 的数据库升级流程与 2.7 及之后相同:
airflow db migrate插件(plugins)注意事项:如果你有使用 Flask-AppBuilder 视图(appbuilder_views)、Flask-AppBuilder 菜单项(appbuilder_menu_items)或 Flask 蓝图(flask_blueprints)的插件,需要把其转换为 FastAPI 应用,或者安装 FAB provider(它为 Airflow 3 提供向后兼容层)。理想情况下,应将插件转换为 Airflow 3 的 Plugin 接口,即外部视图(external_views)、FastAPI 应用(fastapi_apps)与 FastAPI 中间件(fastapi_root_middlewares)。
Helm Chart 用户注意事项:如果使用 Airflow Helm Chart 部署,请对照 Airflow 3 中可用的配置项检查你的 values 配置。所有位于webserver之下的配置项都需要改为apiServer,并且许多参数已被重命名或移除。完整的 Chart 升级清单(values.yaml变更、独立 Dag processor、JWT secret、FAB 默认值、最低 Kubernetes 版本,以及 Chart1.16.0..1.18.0期间重命名的键)参见 chart/docs/upgrading-to-airflow-3.rst。
Step 7:修改启动脚本
在 Airflow 3 中,Webserver 已变成一个通用的 API Server,使用以下命令启动:
airflow api-serverDag Processor 现在必须独立启动,即使是本地或开发环境也是如此:
airflow dag-processor这两个命令都已在 cli_config.py 中注册(分别映射到airflow.cli.commands.api_server_command.api_server与airflow.cli.commands.dag_processor_command.dag_processor)。完成以上步骤后,你应该就能启动 Airflow 3 实例了。
Step 8:升级后需要检查的事项
升级完成后,建议检查以下内容:
- 如果你使用 OAuth、OIDC 或 LDAP 配置了单点登录(SSO),请确认认证正常工作。如果你使用自定义的
webserver_config.py,需要把from airflow.www.security import AirflowSecurityManager替换为from airflow.providers.fab.auth_manager.security_manager.override import FabAirflowSecurityManagerOverride。
破坏性变更(Breaking Changes)完整清单
一些在 Airflow 2.x 中已被弃用的能力在 Airflow 3 中不可用,具体包括:
SubDAGs:被 TaskGroups、Assets 与 Data Aware Scheduling 取代;
SequentialExecutor:被 LocalExecutor 取代,LocalExecutor 可用于 SQLite 的本地开发场景;
CeleryKubernetesExecutor 与 LocalKubernetesExecutor:被 Multiple Executor Configuration(多执行器配置)取代;
SLAs:已弃用并移除,被 Deadline Alerts 取代,参见 deadline-alerts.rst;
Subdir:作为许多 CLI 命令参数的
--subdir或-S已被 Dag Bundles 取代,参见 dag-bundles.rst;REST API(
/api/v1)被取代:请改用基于 FastAPI 的现代化稳定版/api/v2,详见 stable-rest-api-ref.rst;部分 Airflow context 变量被移除:以下键在任务实例的 context 中不再可用。若不替换,将导致 Dag 报错:
tomorrow_dstomorrow_ds_nodashyesterday_dsyesterday_ds_nodashprev_dsprev_ds_nodashprev_execution_dateprev_execution_date_successnext_execution_datenext_ds_nodashnext_dsexecution_date
catchup_by_defaultDag 参数默认值改为False;create_cron_data_intervals配置默认值改为False:这意味着默认将使用CronTriggerTimetable而非CronDataIntervalTimetable。这只影响向
schedule=传递裸 cron 字符串的 Dag(例如schedule="0 0 * * *");传递显式 timetable 实例的 Dag 不受影响。请判断你是否依赖data_interval_start/data_interval_end(以及任务中相关的模板值如ds/ts,它们由logical_date派生,并会在两种 timetable 之间发生偏移)。如果依赖,请显式设置create_cron_data_intervals=True以继续使用CronDataIntervalTimetable;如果不依赖,新的False默认值没有问题。必须在升级前设置该参数。如果改为在已有 Airflow 3 dagruns 之后再修改此标志(从
CronTriggerTimetable切换到CronDataIntervalTimetable),会跳过一次调度运行,以避免与前一次运行的logical_date冲突。手动 Dag 运行与数据间隔(data intervals):在 Airflow 3 中,不要假设手动触发的 Dag run 的
data_interval由(或等于)所提供的logical_date派生。如果 Dag 逻辑需要用户指定的触发日期,请显式使用logical_date。这尤其影响在手动触发时或使用TriggerDagRunOperator时读取data_interval_start或data_interval_end的工作流(详细迁移指导见下文"手动 Dag 运行与logical_date"一节);Simple Auth 现在是默认的
auth_manager:要继续使用 FAB 作为 Auth Manager,请安装 FAB provider 并设置auth_manager为FabAuthManager:airflow.providers.fab.auth_manager.fab_auth_manager.FabAuthManagerAUTH API 路由前缀变化:auth manager 中定义的 api 路由以
/auth路由为前缀。应用外部消费的 URL(如 oauth 重定向 URL)需要相应更新。例如,Airflow 2.x 中的 oauth 重定向 URLhttps://<your-airflow-url.com>/oauth-authorized/google在 Airflow 3.x 中将是https://<your-airflow-url.com>/auth/oauth-authorized/google;XCom Pull 默认行为变化:不带
task_ids参数调用xcom_pull()现在只从当前任务拉取。在 Airflow 2 中,省略task_ids会搜索 Dag run 中的所有任务并返回给定 key 最近推送的值。现在必须显式传入task_ids才能从其他任务拉取 XCom:# Airflow 2 - 从任何任务拉取最近的值 value = ti.xcom_pull(key="shared_state") # Airflow 3 - 同样的调用只检查当前任务 value = ti.xcom_pull(key="shared_state") # Airflow 3 - 指定 task_ids 从其他任务拉取 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派生。最终的data_interval取决于 timetable 与触发路径,某些 API 也允许显式提供数据间隔。
这对以下类型的 Dag 影响最大:
- 在手动运行期间使用
data_interval_start或data_interval_end的 Dag; - 使用
TriggerDagRunOperator触发下游 Dag 的工作流; - 从 Airflow 2 迁移而来、并把
data_interval_start当作手动运行请求日期的 Dag。
迁移指导:如果 Dag 逻辑需要手动运行的用户指定日期,请显式使用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 升级时,请复查所有读取data_interval_start或data_interval_end的手动触发工作流,确认它们真正想要的是间隔语义,还是请求的 logical date。
升级路线小结
把整份指南浓缩为一条可执行的路径:先确认版本与 Python 环境满足要求(Step 1)→ 备份并清理数据库(Step 2)→ 用 ruff + AIR 规则扫描并修复 Dag 导入与破坏性用法(Step 3)→ 安装 Standard Provider(Step 4)→ 消除任务代码中的直接数据库访问(Step 5)→ 用airflow config update --fix迁移配置、用airflow db migrate升级数据库、处理插件与 Helm values(Step 6)→ 将启动脚本切换为airflow api-server+airflow dag-processor(Step 7)→ 最后核对 SSO 认证等收尾事项(Step 8)。其中"任务代码与元数据库解耦"是贯穿始终的主线,理解这一点,Airflow 3 的绝大多数变更都能顺理成章地理解与适配。
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考