Managed Airflow(第 3 代) | Managed Airflow(第 2 代) | Managed Airflow(旧版第 1 代)
本页面介绍如何将 DAG、数据和配置从现有 Airflow 2 的 Managed Airflow(第 3 代)环境转移到 Airflow 3 的 Managed Airflow(第 3 代)环境。
其他迁移指南
| 发件人 | 收件人 | 方法 | 指南 |
|---|---|---|---|
| Managed Airflow(第 3 代),Airflow 2 | Managed Airflow(第 3 代),Airflow 3 | 并列手动转移 | 这一指南 |
| Managed Airflow(第 2 代) | Managed Airflow(第 3 代) | 使用迁移脚本进行并排比较 | 脚本迁移指南 |
| Managed Airflow(第 2 代) | Managed Airflow(第 3 代) | 并排比较(使用快照) | 快照迁移指南 |
| Managed Airflow(旧版第 1 代)、Airflow 2 | Managed Airflow(第 3 代) | 并排比较(使用快照) | 快照迁移指南 |
| Managed Airflow(旧版第 1 代)、Airflow 2 | Managed Airflow(第 2 代) | 并排比较(使用快照) | 快照迁移指南 |
| Managed Airflow(旧版第 1 代)、Airflow 2 | Managed Airflow(第 2 代) | 并列手动转移 | 手动迁移指南 |
| Managed Airflow(旧版第 1 代)、Airflow 1 | Managed Airflow(第 2 代)、Airflow 2 | 并排比较(使用快照) | 快照迁移指南 |
| Managed Airflow(旧版第 1 代)、Airflow 1 | Managed Airflow(第 2 代)、Airflow 2 | 并列手动转移 | 手动迁移指南 |
| Managed Airflow(旧版第 1 代)、Airflow 1 | Managed Airflow(旧版第 1 代)、Airflow 2 | 并列手动转移 | 手动迁移指南 |
Airflow 3 中引入的更改
在开始将 Managed Airflow 环境与 Airflow 3 结合使用之前,请考虑 Airflow 3 为 Managed Airflow(第 3 代)环境所做的更改。
如需大致了解 Airflow 3 社区版中引入的变更,请参阅 Apache Airflow 3 正式版发布!。
DAG 版本控制
在 Airflow 3 中,DAG 将根据开始时的版本运行直至完成,即使在 DAG 运行期间上传了新版本也是如此。
- Airflow 界面中的所有 DAG 运行现在都与相应的 DAG 版本(执行时的版本)相关联。这包括任务结构和 DAG 代码。
回填方面的改进
Airflow 3 对回填(针对历史数据重新运行流水线)的处理方式进行了重大改进。回填从手动流程转变为集成到核心 Airflow 引擎中的完全可观测功能:
- 现在,回填直接在 Airflow 调度程序中进行管理,而不是作为单独的手动流程进行处理。这有助于提高可伸缩性并实现更精确的控制。
- 现在,除了 Airflow CLI 之外,您还可以直接通过 Airflow 界面或 API 调用来触发、停止和监控回填进度。
- Airflow 调度程序可让您更好地了解历史回填运行的状态和健康状况。
- 虽然机器学习社区强烈要求(在旧数据上重新训练模型),但回填改进适用于所有 ETL/ELT 工作流。
提升了安全性和可靠性
在 Airflow 3 中,任务只能通过任务 SDK 与中央 API 服务器通信(在 Airflow 2 中,任务可以直接访问数据库)。API 服务器会高效地汇集这些连接。您的数据库可免受连接峰值的影响,从而使整个环境在重负荷下更加稳定。
通过利用新的任务执行接口,Airflow 3 支持更好地隔离任务,防止一个任务可能干扰或访问另一个任务的数据。
Airflow 3 的 CLI 不再直接访问数据库。新的
airflowctl命令行界面是一个单独的软件包,专门用于通过 API 进行远程访问。它不会直接访问数据库,而是通过 API 与 Airflow 交互,这样更安全。
事件驱动型调度和数据资产
数据集已演变为数据资产。借助数据资产,Airflow 可以更好地跟踪和响应由 Airflow 外部的系统创建或更新的数据。
Airflow 3 中引入了一个名为 Watcher 的新概念。这些组件会监控数据资产的变化,以便 Airflow 在数据到达时立即触发工作流。现在,DAG 不再需要轮询(每分钟检查一次文件是否存在),而是可以在消息到达消息队列的那一刻立即触发。
Airflow 3 引入了使用 Python 装饰器的新资产中心语法,使代码更清晰,更直观。
经过现代化改造的 Airflow 界面
- Airflow 界面使用 React(前端)和 FastAPI(后端)从头开始重写。
- 新的 Airflow 界面通过标准化的 REST API 和专门用于界面操作的 API 来执行操作。
- 通过将 Flask 实现替换为 FastAPI,Airflow 界面变得更加灵敏。
- 网格视图和图表视图已统一,可实现更顺畅的工作流程,让您更轻松地在高层 DAG 结构和特定任务日志之间切换。
Airflow 3 中的重大变更
Airflow 3 引入了一些重大更改,其中一些更改重大:
- Airflow 2 中的现有 DAG 无法保证与 Airflow 3 搭配使用。必须通过更改导入、DAG 参数和其他实现细节来测试并可能调整它们。
某些 Airflow 2 配置选项在 Airflow 3 中已重命名或移除。如需详细了解参数,请参阅 Airflow 配置参考文档。
无法从任务代码直接访问 Airflow 数据库:
- 任务代码无法再直接导入和使用 Airflow 数据库会话或模型。
- 无法将
PostgresHook和PostgresOperator与airflow_db连接搭配使用。
某些自定义 PyPI 软件包可能与新版 Airflow 及其依赖项不兼容。
REST API (
/api/v1) 已替换为/api/v2。SubDAG 已被 TaskGroup、Asset 和 Data Aware Scheduling 取代。
我们已弃用并移除 SLA。这些提醒已替换为截止日期提醒。
移除了 CLI 命令中的 subdir 实参。
部分 Airflow 上下文变量已被移除。如需了解详情,请参阅 Airflow 文档中的重大更改。
catchup_by_defaultDAG 参数现在默认为False。create_cron_data_intervals配置现在默认处于False状态。 这意味着默认情况下将使用CronTriggerTimetable,而不是CronDataIntervalTimetable。
Airflow 3 和 Airflow 2 的环境之间存在差异
使用 Airflow 2 的 Managed Airflow 环境与使用 Airflow 3 的环境之间的主要区别如下:
Airflow 3 环境中的工作负载配置:
更改了
[celery]worker_concurrency配置选项的自动计算方式,以适应 Airflow 3 各组件的不同内存用量。在 Airflow 3 中,无法直接从任务代码访问 Airflow 数据库。
Airflow 3 使用
airflowctl命令行实用程序来运行 Airflow CLI 命令。在 Airflow 3 环境中,预安装的 PyPI 软件包有所不同。如需查看预安装的 PyPI 软件包的列表,请参阅预安装的软件包的更改日志。
并行迁移到 Airflow 3
并行迁移过程包含以下步骤:
- 检查与 Airflow 3 的兼容性。
- 创建 Airflow 3 环境,转移配置替换和环境变量。
- 将 PyPI 软件包安装到 Airflow 3 环境。
- 将变量、连接和池转移到 Airflow 3。
- 从 Airflow 2.* 环境存储桶转移其他数据。
- 转移用户和角色。
- 确保您的 DAG 已准备好使用 Airflow 3。
- 将 DAG 转移到 Airflow 3 环境。
- 监控 Airflow 3 环境。
第 1 步:检查与 Airflow 3 的兼容性
如需检查与 Airflow 3 的兼容性,请执行以下操作:
- 检查您的环境是否使用 Airflow 2.7 版或更高版本。我们建议先升级到最新的 Airflow 2 版本,然后再迁移到 Airflow 3。
- 检查环境是否健康,并且已正常运行一段时间。
- 确保您的 DAG 和 Airflow 配置不使用任何已在 Airflow 3 中移除的功能。
- 请参阅有关更改 DAG 以使其与 Airflow 3 兼容的说明,了解在迁移过程中是否需要对 DAG 进行任何更改。
- 使用社区版 Airflow 提供的
ruff工具检查 Airflow DAG 的兼容性。如需查看相关说明,请参阅 Airflow 文档中的检查 Airflow DAG 的兼容性。
第 2 步:创建 Airflow 3 环境,转移配置替换和环境变量
在此步骤中,您将创建一个新的 Managed Airflow(第 3 代)环境,其中包含 Airflow 3,并开始从 Airflow 2 环境转移配置参数:
按照创建 Managed Airflow (Gen 3) 环境的步骤执行以下操作:
- 选择 Airflow build 时,请选择具有 Airflow 3 的 build。
从 Airflow 2 环境中复制所有兼容的 Airflow 配置选项替换。
从 Airflow 2 环境中复制所有环境变量。
继续创建使用 Airflow 3 的环境。
下表列出了一些 Airflow 配置选项更改。该列表并不详尽。如需详细了解 Airflow 配置选项的更改,请参阅 Airflow 文档中的 Airflow 配置参考和 Airflow 版本说明。
| Airflow 2 选项 | Airflow 3 选项 |
|---|---|
[scheduler]min_file_process_interval
|
[dag_processor]min_file_process_interval
|
[webserver]rbac_user_registration_role
|
[api]rbac_user_registration_role
|
[core]dag_file_processor_timeout
|
[dag_processor]dag_file_processor_timeout
|
[scheduler]dag_dir_list_interval
|
[dag_processor]refresh_interval
|
[scheduler]max_threads
|
[dag_processor]parsing_processes
|
[scheduler]parsing_processes
|
[dag_processor]parsing_processes
|
[webserver]instance_name
|
[api]instance_name
|
[scheduler]scheduler_zombie_task_threshold
|
[scheduler]task_instance_heartbeat_timeout
|
[webserver]rbac
|
已弃用 |
[api]auth_backend=airflow.api.auth.backend.deny_all
|
已弃用 |
[api]auth_backends=airflow.api.auth.backend.deny_all
|
已弃用 |
[api]composer_auth_user_registration_role
|
已弃用 |
第 3 步:将 PyPI 软件包安装到 Airflow 3 环境
创建 Airflow 3 环境后,请安装 PyPI 软件包:
- 从 Airflow 2 环境中复制 PyPI 软件包要求。
- 启动 PyPI 软件包更新操作,并等待环境更新。
由于 Airflow 3 环境使用了一组不同的预装软件包,因此您在更新操作期间可能会遇到 PyPI 软件包冲突。如需详细了解如何排查 PyPI 软件包冲突,请参阅与预安装的 PyPI 软件包冲突。
第 4 步:从 Airflow 2 中导出变量、连接和池
如果您没有变量或连接,请跳过相应的导出和导入命令。
您只有在拥有 default_pool 以外的自定义池时才需要转移池。否则,请跳过导出和导入池的命令。
从 Airflow 2 环境中导出变量:
gcloud composer environments run AIRFLOW_2_ENV \ --location AIRFLOW_2_LOCATION \ variables -- export /home/airflow/gcs/data/variables.json替换以下内容:
AIRFLOW_2_ENV:Airflow 2 环境的名称。AIRFLOW_2_LOCATION:Airflow 2 环境所在的区域。
从 Airflow 2 环境中导出连接:
gcloud composer environments run AIRFLOW_2_ENV \ --location AIRFLOW_2_LOCATION \ connections -- export /home/airflow/gcs/data/connections.json替换以下内容:
AIRFLOW_2_ENV:Airflow 2 环境的名称。AIRFLOW_2_LOCATION:Airflow 2 环境所在的区域。
从 Airflow 2 环境中导出池:
gcloud composer environments run AIRFLOW_2_ENV \ --location AIRFLOW_2_LOCATION \ pools -- export /home/airflow/gcs/data/pools.json替换以下内容:
AIRFLOW_2_ENV:Airflow 2 环境的名称。AIRFLOW_2_LOCATION:Airflow 2 环境所在的区域。
获取 Airflow 2 环境的存储桶名称:
gcloud composer environments describe AIRFLOW_2_ENV \ --location AIRFLOW_2_LOCATION \ --format="value(storageConfig.bucket)"替换以下内容:
AIRFLOW_2_ENV:Airflow 2 环境的名称。AIRFLOW_2_LOCATION:Airflow 2 环境所在的区域。
从 Airflow 2 环境的存储桶的
/data目录中下载variables.json、connections.json和pools.json文件到本地目录:gcloud storage cp gs://AIRFLOW_2_BUCKET/data/variables.json ./variables.json gcloud storage cp gs://AIRFLOW_2_BUCKET/data/connections.json ./connections.json gcloud storage cp gs://AIRFLOW_2_BUCKET/data/pools.json ./pools.json替换以下内容:
AIRFLOW_2_BUCKET:您的 Airflow 2 环境的存储桶名称,在上一步中获得。
第 5 步:将变量、连接和池导入 Airflow 3
如果您没有变量或连接,请跳过相应的导出和导入命令。
您只有在拥有 default_pool 以外的自定义池时才需要转移池。否则,请跳过导出和导入池的命令。
配置
airflowctl以针对 Airflow 3 环境运行 Airflow CLI 命令。使用
airflowctl将变量、连接和池导入 Airflow 3 环境:airflowctl variables import ./variables.json airflowctl connections import ./connections.json airflowctl pools import ./pools.json确认变量、连接和池已导入到 Airflow 3 环境:
airflowctl variables list airflowctl connections list airflowctl pools list清理 JSON 文件:
gcloud storage rm gs://AIRFLOW_2_BUCKET/data/variables.json gcloud storage rm gs://AIRFLOW_2_BUCKET/data/connections.json gcloud storage rm gs://AIRFLOW_2_BUCKET/data/pools.json rm ./variables.json rm ./connections.json rm ./pools.json替换以下内容:
AIRFLOW_2_BUCKET:Airflow 2 环境的存储桶名称。
第 6 步:从 Airflow 2 环境的存储桶转移其他数据
在此步骤中,您将从 Airflow 2 环境的存储桶转移剩余数据。
获取 Airflow 3 环境的存储桶名称:
gcloud composer environments describe AIRFLOW_3_ENV \ --location AIRFLOW_3_LOCATION \ --format="value(storageConfig.bucket)"替换以下内容:
AIRFLOW_3_ENV:Airflow 3 环境的名称。AIRFLOW_3_LOCATION:Airflow 3 环境所在的区域。
将插件从 Airflow 2 环境的存储桶导出到 Airflow 3 环境的存储桶中的
/plugins目录:gcloud composer environments storage plugins export \ --destination=AIRFLOW_3_BUCKET/plugins \ --environment=AIRFLOW_2_ENV \ --location=AIRFLOW_2_LOCATION替换以下内容:
AIRFLOW_3_BUCKET:Airflow 3 环境的存储桶名称,在上一步中获取。AIRFLOW_2_ENV:Airflow 2 环境的名称。AIRFLOW_2_LOCATION:Airflow 2 环境所在的区域。
检查
/plugins目录是否已成功导入:gcloud composer environments storage plugins list \ --environment=AIRFLOW_3_ENV \ --location=AIRFLOW_3_LOCATION替换以下内容:
AIRFLOW_3_ENV:Airflow 3 环境的名称。AIRFLOW_3_LOCATION:Airflow 3 环境所在的区域。
将
/data目录从 Airflow 2 环境导出到 Airflow 3 环境:gcloud composer environments storage data export \ --destination=AIRFLOW_3_BUCKET/data \ --environment=AIRFLOW_2_ENV \ --location=AIRFLOW_2_LOCATION替换以下内容:
AIRFLOW_3_BUCKET:Airflow 3 环境的存储桶名称,在上一步中获取。AIRFLOW_2_ENV:Airflow 2 环境的名称。AIRFLOW_2_LOCATION:Airflow 2 环境所在的区域。
检查
/data文件夹是否已成功导入:gcloud composer environments storage data list \ --environment=AIRFLOW_3_ENV \ --location=AIRFLOW_3_LOCATION
第 7 步:转移用户和角色
无法迁移用户和角色,因为 airflowctl 尚不支持 users 和 roles 命令。
第 8 步:确保您的 DAG 已准备好使用 Airflow 3
调整 Airflow DAG,使其与 Airflow 3 兼容。
查看直接访问 Airflow 数据库的自定义编写的任务:
在 Airflow 3 中,运算符无法使用数据库会话直接访问 Airflow 元数据数据库。如果您有自定义运算符,请检查您的代码,确保没有直接的数据库访问调用。
您可以采用以下替代方法之一,在任务中迁移到不直接访问 Airflow 数据库:
通过将 Airflow 数据库内容导出到 Cloud SQL 实例来访问 Airflow 数据库。
使用 Airflow Python 客户端。 Airflow 社区版提供的 Python 客户端针对大多数表(例如
DagRuns、TaskInstances、Variables、Connections、XComs)定义了 API。apache-airflow-client软件包已预安装在 Managed Airflow Airflow 3 build 中。通过 DAG 运行
airflowctl到BashOperator。
如果查询导出的 Airflow 数据库不适合您的使用情形,并且 Airflow Python 客户端和
airflowctl均未提供所需的功能,请考虑在 Airflow 的社区版中请求新的 API 端点或任务 SDK 功能。如果您有
KubernetesExecutor个任务,请将queue="kubernetes"替换为executor="KubernetesExecutor",以调整其运算符定义。Airflow 3 中
KubernetesExecutor任务的示例:PythonOperator( task_id="airflow3_kubernetes_executor_task", dag=dag, python_callable=f, executor="KubernetesExecutor", )如果您在任务代码中使用
AIRFLOW__WEBSERVER__BASE_URL环境变量,请将其替换为[api]base_urlAirflow 配置选项。在 Airflow 3 中获取此值的示例:
from airflow.configuration import conf webserver_base_url = conf.get("api", "base_url")
第 9 步:将 DAG 转移到 Airflow 3 环境
在环境之间转移 DAG 时,可能会出现以下问题:
如果两个环境中都启用了(未暂停)某个 DAG,则每个环境都会安排自己的 DAG 副本。这可能会导致在相同的数据和执行时间内并发运行 DAG。
由于 DAG 同步,Airflow 从 DAG 中指定的开始日期开始安排额外的 DAG 运行。这是因为新的 Airflow 实例不考虑从 Airflow 2 环境运行的 DAG 的历史记录。这可能会导致从指定的开始日期开始安排大量 DAG 运行。
防止并发 DAG 运行
在您的 Airflow 3 环境中,替换 dags_are_paused_at_creation Airflow 配置选项。进行此项更改后,默认情况下,所有新 DAG 都会暂停。
| 部分 | 键 | 值 |
|---|---|---|
core |
dags_are_paused_at_creation |
True |
防止额外运行或缺失的 DAG 运行
在要转移到 Airflow 3 环境的 DAG 中,指定新的静态开始日期。
为避免逻辑日期出现间断和重叠,第一次 DAG 运行必须在下一个安排时间间隔的 Airflow 3 环境中进行。为此,请将 DAG 中新的开始日期设置为早于 Airflow 2 环境中最后一次运行的日期。
例如,如果您在 Airflow 2 环境中的每天 15:00、17:00 和 21:00 运行 DAG,则上次 DAG 运行发生在 15:00,并且您计划于 15:15 转移此 DAG,那么 Airflow 3 环境的开始日期就是今天的 14:45。在 Airflow 3 环境中启用 DAG 后,Airflow 会安排在 17:00 运行 DAG。
再举一个例子,如果您的 DAG 在 Airflow 2 环境中的每一天 00:00 运行,则上次 DAG 运行发生在 2026 年 3 月 26 日 00:00,并且您计划在 2026 年 3 月 26 日 13:00 转移该 DAG,Airflow 3 环境的开始日期可以是 2026 年 3 月 25 日 23:45。在 Airflow 3 环境中启用 DAG 后,Airflow 将安排在 2026 年 3 月 27 日 00:00 运行 DAG。
将 DAG 逐个转移到 Airflow 3 环境
对于每个 DAG,请按照以下步骤转移它:
确保按照上一部分中的说明设置 DAG 中的新开始日期。
将更新的 DAG 上传到 Airflow 3 环境。由于配置替换,此 DAG 在 Airflow 3 环境中暂停,因此尚未安排任何 DAG 运行。
在 Airflow 网页界面中,前往 DAG 并查看报告的 DAG 语法错误。
当您计划转移 DAG 时,请执行以下操作:
暂停 Airflow 2 环境中的 DAG。
取消暂停 Airflow 3 环境中的 DAG。
检查新的 DAG 运行是否安排在正确的时间。
等待 DAG 在 Airflow 3 环境中运行,并检查运行是否成功。
根据 DAG 是否成功运行:
如果 DAG 运行成功,您可以继续在 Airflow 3 环境中使用 DAG。最后,考虑删除 Airflow 2 版本的 DAG。
如果 DAG 运行失败,请尝试排查 DAG 问题,直到它在 Airflow 3 中成功运行为止。
如果需要,您始终可以回退到 Airflow 2 版本的 DAG:
暂停 Airflow 3 环境中的 DAG。
在 Airflow 3 环境中取消暂停 DAG。这会为新的 DAG 运行安排与失败的 DAG 运行相同的日期和时间。
当您准备好继续使用 Airflow 3 版本的 DAG 时,请调整开始日期,将新版本的 DAG 上传到 Airflow 3 环境,然后重复上述过程。
第 10 步:监控 Airflow 3 环境
将所有 DAG 和配置转移到 Airflow 3 环境后,请监控其是否存在潜在问题、DAG 运行失败以及环境整体运行状况。如果 Airflow 3 环境在足够长的时间内正常运行,则可以移除 Airflow 2 环境。