将环境迁移到 Airflow 3(并行迁移)

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 数据库会话或模型。
    • 无法将 PostgresHookPostgresOperatorairflow_db 连接搭配使用。
  • 某些自定义 PyPI 软件包可能与新版 Airflow 及其依赖项不兼容。

  • REST API (/api/v1) 已替换为 /api/v2

  • SubDAG 已被 TaskGroup、Asset 和 Data Aware Scheduling 取代。

  • 我们已弃用并移除 SLA。这些提醒已替换为截止日期提醒

  • 移除了 CLI 命令中的 subdir 实参。

  • 部分 Airflow 上下文变量已被移除。如需了解详情,请参阅 Airflow 文档中的重大更改

  • catchup_by_default DAG 参数现在默认为 False

  • create_cron_data_intervals 配置现在默认处于 False 状态。 这意味着默认情况下将使用 CronTriggerTimetable,而不是 CronDataIntervalTimetable

  • Airflow 3.0.0 的变更列表

  • Airflow 3.1.0 的变更列表

Airflow 3 和 Airflow 2 的环境之间存在差异

使用 Airflow 2 的 Managed Airflow 环境与使用 Airflow 3 的环境之间的主要区别如下:

  • Airflow 3 环境中的工作负载配置:

    • 所有 Airflow 组件的最低内存量为 2 GB。
    • 与 Airflow 2 环境相比,环境预设中 Airflow 触发器和工作器的配置发生了变化。
    • 触发器的默认 CPU 数量为 1。
    • 触发器的默认内存量为 2 GB。
  • 更改了 [celery]worker_concurrency 配置选项的自动计算方式,以适应 Airflow 3 各组件的不同内存用量。

  • 在 Airflow 3 中,无法直接从任务代码访问 Airflow 数据库

  • Airflow 3 使用 airflowctl 命令行实用程序来运行 Airflow CLI 命令。

  • 在 Airflow 3 环境中,预安装的 PyPI 软件包有所不同。如需查看预安装的 PyPI 软件包的列表,请参阅预安装的软件包的更改日志

并行迁移到 Airflow 3

并行迁移过程包含以下步骤:

  1. 检查与 Airflow 3 的兼容性。
  2. 创建 Airflow 3 环境,转移配置替换和环境变量。
  3. 将 PyPI 软件包安装到 Airflow 3 环境。
  4. 将变量、连接和池转移到 Airflow 3。
  5. 从 Airflow 2.* 环境存储桶转移其他数据。
  6. 转移用户和角色。
  7. 确保您的 DAG 已准备好使用 Airflow 3。
  8. 将 DAG 转移到 Airflow 3 环境。
  9. 监控 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) 环境的步骤执行以下操作:

  1. 选择 Airflow build 时,请选择具有 Airflow 3 的 build。
  2. 从 Airflow 2 环境中复制所有兼容的 Airflow 配置选项替换。

  3. 从 Airflow 2 环境中复制所有环境变量。

  4. 继续创建使用 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 软件包

  1. 从 Airflow 2 环境中复制 PyPI 软件包要求。
  2. 启动 PyPI 软件包更新操作,并等待环境更新。

由于 Airflow 3 环境使用了一组不同的预装软件包,因此您在更新操作期间可能会遇到 PyPI 软件包冲突。如需详细了解如何排查 PyPI 软件包冲突,请参阅与预安装的 PyPI 软件包冲突

第 4 步:从 Airflow 2 中导出变量、连接和池

如果您没有变量或连接,请跳过相应的导出和导入命令。

您只有在拥有 default_pool 以外的自定义池时才需要转移池。否则,请跳过导出和导入池的命令。

  1. 从 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 环境所在的区域。
  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 环境所在的区域。
  3. 从 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 环境所在的区域。
  4. 获取 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 环境所在的区域。
  5. 从 Airflow 2 环境的存储桶的 /data 目录中下载 variables.jsonconnections.jsonpools.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 以外的自定义池时才需要转移池。否则,请跳过导出和导入池的命令。

  1. 配置 airflowctl 以针对 Airflow 3 环境运行 Airflow CLI 命令。

  2. 使用 airflowctl 将变量、连接和池导入 Airflow 3 环境:

    airflowctl variables import ./variables.json
    airflowctl connections import ./connections.json
    airflowctl pools import ./pools.json
    
  3. 确认变量、连接和池已导入到 Airflow 3 环境:

    airflowctl variables list
    airflowctl connections list
    airflowctl pools list
    
  4. 清理 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 环境的存储桶转移剩余数据。

  1. 获取 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 环境所在的区域。
  2. 将插件从 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 环境所在的区域。
  3. 检查 /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 环境所在的区域。
  4. /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 环境所在的区域。
  5. 检查 /data 文件夹是否已成功导入:

    gcloud composer environments storage data list \
      --environment=AIRFLOW_3_ENV \
      --location=AIRFLOW_3_LOCATION
    

第 7 步:转移用户和角色

无法迁移用户和角色,因为 airflowctl 尚不支持 usersroles 命令。

第 8 步:确保您的 DAG 已准备好使用 Airflow 3

  1. 调整 Airflow DAG,使其与 Airflow 3 兼容

  2. 查看直接访问 Airflow 数据库的自定义编写的任务:

    在 Airflow 3 中,运算符无法使用数据库会话直接访问 Airflow 元数据数据库。如果您有自定义运算符,请检查您的代码,确保没有直接的数据库访问调用。

    您可以采用以下替代方法之一,在任务中迁移到不直接访问 Airflow 数据库:

    • 通过将 Airflow 数据库内容导出到 Cloud SQL 实例来访问 Airflow 数据库。

    • 使用 Airflow Python 客户端。 Airflow 社区版提供的 Python 客户端针对大多数表(例如 DagRunsTaskInstancesVariablesConnectionsXComs)定义了 API。apache-airflow-client 软件包已预安装在 Managed Airflow Airflow 3 build 中。

    • 通过 DAG 运行 airflowctlBashOperator

    如果查询导出的 Airflow 数据库不适合您的使用情形,并且 Airflow Python 客户端和 airflowctl 均未提供所需的功能,请考虑在 Airflow 的社区版中请求新的 API 端点或任务 SDK 功能。

  3. 如果您有 KubernetesExecutor 个任务,请将 queue="kubernetes" 替换为 executor="KubernetesExecutor",以调整其运算符定义。

    Airflow 3 中 KubernetesExecutor 任务的示例:

    PythonOperator(
    task_id="airflow3_kubernetes_executor_task",
    dag=dag,
    python_callable=f,
    executor="KubernetesExecutor",
    )
    
  4. 如果您在任务代码中使用 AIRFLOW__WEBSERVER__BASE_URL 环境变量,请将其替换为 [api]base_url Airflow 配置选项。

    在 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,请按照以下步骤转移它:

  1. 确保按照上一部分中的说明设置 DAG 中的新开始日期。

  2. 将更新的 DAG 上传到 Airflow 3 环境。由于配置替换,此 DAG 在 Airflow 3 环境中暂停,因此尚未安排任何 DAG 运行。

  3. Airflow 网页界面中,前往 DAG 并查看报告的 DAG 语法错误。

  4. 当您计划转移 DAG 时,请执行以下操作:

    1. 暂停 Airflow 2 环境中的 DAG。

    2. 取消暂停 Airflow 3 环境中的 DAG。

    3. 检查新的 DAG 运行是否安排在正确的时间。

    4. 等待 DAG 在 Airflow 3 环境中运行,并检查运行是否成功。

  5. 根据 DAG 是否成功运行:

    • 如果 DAG 运行成功,您可以继续在 Airflow 3 环境中使用 DAG。最后,考虑删除 Airflow 2 版本的 DAG。

    • 如果 DAG 运行失败,请尝试排查 DAG 问题,直到它在 Airflow 3 中成功运行为止。

      如果需要,您始终可以回退到 Airflow 2 版本的 DAG:

      1. 暂停 Airflow 3 环境中的 DAG。

      2. 在 Airflow 3 环境中取消暂停 DAG。这会为新的 DAG 运行安排与失败的 DAG 运行相同的日期和时间。

      3. 当您准备好继续使用 Airflow 3 版本的 DAG 时,请调整开始日期,将新版本的 DAG 上传到 Airflow 3 环境,然后重复上述过程。

第 10 步:监控 Airflow 3 环境

将所有 DAG 和配置转移到 Airflow 3 环境后,请监控其是否存在潜在问题、DAG 运行失败以及环境整体运行状况。如果 Airflow 3 环境在足够长的时间内正常运行,则可以移除 Airflow 2 环境。

后续步骤