十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Apache Airflow 3 升级全指南:从 2.x 迁移的架构变化、破坏性变更与分步实操

Apache Airflow 3 升级全指南:从 2.x 迁移的架构变化、破坏性变更与分步实操 Apache Airflow 3 升级全指南从 2.x 迁移的架构变化、破坏性变更与分步实操【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflowApache 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 ServerAPI Server 目前是任务tasks和 Worker 访问元数据库的唯一入口承载了多种应用形态Airflow REST API、为 Airflow UI托管静态 JS服务的内部 API、以及 Worker 通过任务执行接口task execution interface执行 TaskInstanceTI时交互所用的 APIWorker 不再直连数据库而是与 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 3Python 版本必须在受支持列表中先确认你的运行环境满足 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 3Dag 作者——检查 Dag 的兼容性为了最小化升级摩擦Airflow 社区基于Ruff与AIRAirflow规则创建了 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.sdkairflow.decorators.dagairflow.sdk.dagairflow.decorators.taskairflow.sdk.taskairflow.decorators.task_groupairflow.sdk.task_groupairflow.decorators.setupairflow.sdk.setupairflow.decorators.teardownairflow.sdk.teardownairflow.models.dag.DAGairflow.sdk.DAGairflow.models.baseoperator.BaseOperatorairflow.sdk.BaseOperatorairflow.models.param.Paramairflow.sdk.Paramairflow.models.param.ParamsDictairflow.sdk.ParamsDictairflow.models.baseoperatorlink.BaseOperatorLinkairflow.sdk.BaseOperatorLinkairflow.sensors.base.BaseSensorOperatorairflow.sdk.BaseSensorOperatorairflow.hooks.base.BaseHookairflow.sdk.BaseHookairflow.notifications.basenotifier.BaseNotifierairflow.sdk.BaseNotifierairflow.utils.task_group.TaskGroupairflow.sdk.TaskGroupairflow.utils.context.Contextairflow.sdk.Contextairflow.datasets.Datasetairflow.sdk.Assetairflow.datasets.DatasetAliasairflow.sdk.AssetAliasairflow.datasets.DatasetAllairflow.sdk.AssetAllairflow.datasets.DatasetAnyairflow.sdk.AssetAnyairflow.models.connection.Connectionairflow.sdk.Connectionairflow.models.variable.Variableairflow.sdk.Variableairflow.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 能力而非开放直接数据库访问。已知的变通方案使用 DbApiHookPostgresHook 或 MySqlHook警告此方案不被推荐仅作为无法使用 Python Client 的用户的已知变通方案被记录。该方案存在显著限制并且在未来的 Airflow 版本中会失效。需要重点考虑的事项未来版本会失效该方案在 Airflow 3.2 及之后会失效schema 变更时你需要自行负责适配代码数据库 schema 不是公共 APIAirflow 元数据库 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_idmetadata_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也可以使用SQLExecuteQueryOperatorfrom airflow.providers.common.sql.operators.sql import SQLExecuteQueryOperator query_task SQLExecuteQueryOperator( task_idquery_metadata, conn_idmetadata_postgres, sqlSELECT conn_id, conn_type FROM connection WHERE conn_type postgres, do_xcom_pushTrue, )注意为元数据库连接始终使用只读数据库凭据并建议使用临时凭据。Step 6部署管理人员——升级 Airflow 实例为让升级过程更简单、更安全Airflow 3 提供了配置升级工具。先从配置检查开始airflow config update该工具也能把配置自动更新为兼容 Airflow 3airflow config update --fix从源码看该命令的实现位于 config_command.pylint_config函数扫描airflow.cfg中在 Airflow 3.0 里被移除或重命名的参数并给出建议update_config函数默认执行 dry-run仅展示变更--fix才会真正改写airflow.cfg改写前会先备份旧文件且会清理原有注释--all-recommendations则把非破坏性的推荐变更一并纳入。命令内置的CONFIGS_CHANGES列表覆盖了数百项迁移规则例如core.executor默认值由SequentialExecutor改为LocalExecutorcore.sql_alchemy_conn等一批参数从core段迁移到database段scheduler.catchup_by_default默认值由True改为Falsescheduler.create_cron_data_intervals与create_delta_data_intervals默认值均由True改为Falsewebserver段的web_server_host/web_server_port等参数迁移到api段的host/portscheduler.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.rstSubdir作为许多 CLI 命令参数的--subdir或-S已被 Dag Bundles 取代参见 dag-bundles.rstREST 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_datecatchup_by_defaultDag 参数默认值改为Falsecreate_cron_data_intervals配置默认值改为False这意味着默认将使用CronTriggerTimetable而非CronDataIntervalTimetable。这只影响向schedule传递裸 cron 字符串的 Dag例如schedule0 0 * * *传递显式 timetable 实例的 Dag 不受影响。请判断你是否依赖data_interval_start/data_interval_end以及任务中相关的模板值如ds/ts它们由logical_date派生并会在两种 timetable 之间发生偏移。如果依赖请显式设置create_cron_data_intervalsTrue以继续使用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为FabAuthManagerairflow.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/googleXCom Pull 默认行为变化不带task_ids参数调用xcom_pull()现在只从当前任务拉取。在 Airflow 2 中省略task_ids会搜索 Dag run 中的所有任务并返回给定 key 最近推送的值。现在必须显式传入task_ids才能从其他任务拉取 XCom# Airflow 2 - 从任何任务拉取最近的值 value ti.xcom_pull(keyshared_state) # Airflow 3 - 同样的调用只检查当前任务 value ti.xcom_pull(keyshared_state) # Airflow 3 - 指定 task_ids 从其他任务拉取 value ti.xcom_pull(task_idsupstream_task, keyshared_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_datefrom airflow.decorators import get_current_context, task task def process_data(): context get_current_context() processing_date context[logical_date] return fProcessing 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 ProviderStep 4→ 消除任务代码中的直接数据库访问Step 5→ 用airflow config update --fix迁移配置、用airflow db migrate升级数据库、处理插件与 Helm valuesStep 6→ 将启动脚本切换为airflow api-serverairflow dag-processorStep 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),仅供参考
返回列表