资讯动态

Apache Airflow 插件系统完全指南:从 plugins 目录到 FastAPI、React App 与视图扩展

发布时间:2026/9/9 20:53:53 来源:尧图企业网站定制
Apache Airflow 插件系统完全指南从 plugins 目录到 FastAPI、React App 与视图扩展【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflowAirflow 内置了一套轻量级的插件管理器通过把 Python 文件放入$AIRFLOW_HOME/plugins目录就能把外部功能宏、自定义视图、API 端点、自定义调度逻辑等注入核心系统。本文以 Airflow 官方管理文档 为主干结合 Airflow 3.x 源码系统讲解插件的加载时机、AirflowPlugin编程接口、外部视图External Views、React App、FastAPI App 与中间件、applies_to作用域控制、entrypoint 打包分发以及故障排查方法帮助你在不修改 Airflow 核心的前提下深度定制自己的平台。插件是什么为什么需要它Airflow 提供的是通用的数据处理工具箱而不同组织的数据栈和需求各不相同。通过插件机制公司可以把 Airflow 安装定制成贴合自身生态的平台实现「编写、共享、激活」新功能集的轻量途径同时也满足了复杂应用与各种数据、元数据交互的需求。官方文档给出的典型插件应用场景包括解析 Hive 日志并暴露 Hive 元数据CPU / IO / 阶段 / 倾斜等的工具集允许收集指标、设置阈值与告警的异常检测框架帮助理解谁访问了什么内容的审计工具配置驱动的 SLA 监控工具监控指定表在什么时间落地、发出告警并可视化故障。选择在 Airflow 之上做二次开发是因为它能直接复用大量现成组件可渲染视图的 Web 服务器、存储模型的元数据库、既有的数据库连接能力与连接知识、可推送负载的 worker 阵列、已经被部署好的基础设施以及底层的图表库与抽象能力。插件可注册的构建模块Building Blocks一个插件可以向 Airflow 注册如下组件External Views外部视图在 UI 中新增按钮 / 标签页链接到新页面。React Apps在 Airflow UI 中嵌入自定义 React 应用Airflow 3.1 新增目前标记为实验性。FastAPI Apps添加自定义 API 端点。FastAPI Middlewares拦截并修改 API 请求 / 响应。Macros宏定义可在 DAG 模板中复用的 Python 函数。Operator Extra Links在任务详情视图中添加自定义按钮。Timetables Listeners实现自定义调度逻辑与事件钩子。Deadline References注册自定义的 Deadline Alert 引用类供 DAG 作为自定义截止时间使用。关于 Operator Extra Links 更深入的开发指引可参考 如何定义 Extra Link。插件的加载机制与生效时机插件从哪来插件可以来自两个渠道见 加载入口 的_get_plugins()实现settings.PLUGINS_FOLDER即$AIRFLOW_HOME/plugins对应配置项[core] plugins_folder默认值为{AIRFLOW_HOME}/plugins目录下的所有 Python 文件已安装包中通过airflow.pluginsentrypoint 注册的插件_load_entrypoint_plugins()Provider 包自带的插件在关闭懒加载时通过_load_providers_plugins()一并载入。加载时管理器会逐个实例化插件并调用其on_load()用loaded_plugins集合做按名称去重——若两个插件同名后加载者会被跳过并记录Plugin name already registered, skipping警告源码位于 plugins_manager.py。懒加载与重载时机默认情况下插件是懒加载的一旦加载便不会重载唯一例外是 UI 插件会在 Web 服务器中自动加载。因此修改插件后想让 webserver、scheduler 使用新代码必须重启相应进程新代码在调度器scheduler启动前不会反映到正在运行的新任务中。若希望在每次启动 Airflow 进程时都强制加载插件可在airflow.cfg中设置[core] lazy_load_plugins False该项对应的默认值为True定义见 config.yml。fork 执行对插件更新的影响任务执行默认使用fork进程派生以避免每次创建新 Python 解释器、重新解析 Airflow 全部代码与启动例程带来的开销——这对短任务尤其受益。但这也意味着如果任务代码里使用了插件想让插件更新生效要么重启 workerCeleryExecutor 下或 schedulerLocalExecutor 下要么接受启动时的性能损失把core.execute_tasks_new_python_interpreter设为True让每个任务都启动全新的 Python 解释器。该配置项的定义见 config.yml默认False采用 fork设为True后启动速度较慢但插件改动能立即被新任务拾取。反过来仅被 DAG 文件导入的模块不受此问题困扰因为 DAG 文件不会在常驻的 Airflow 进程中被加载 / 解析。使用airflow plugins命令排障当插件出现问题时可以用airflow plugins命令输出已加载插件的详细信息命令行实现位于 plugins_command.py是排查「插件到底加载了没有」的第一步。AirflowPlugin 编程接口创建插件需要继承airflow.plugins_manager.AirflowPlugin并在类中引用想要挂接进 Airflow 的对象。基类定义如下完整字段class AirflowPlugin: # 插件名称 (str) name None # 注入宏命名空间的引用列表 macros [] # 包含 FastAPI app 对象及元数据的字典列表见下文示例 fastapi_apps [] # 包含 FastAPI 中间件工厂对象及元数据的字典列表见下文示例 fastapi_root_middlewares [] # 包含外部视图及元数据的字典列表见下文示例 external_views [] # 包含 react app 及元数据的字典列表见下文示例 # 注意React app 仅 Airflow 3.1 及以上支持集成目前为实验性接口未来可能变化 # 尤其是较复杂插件应用中 UI 与插件的依赖、状态交互可能需要重构。 react_apps [] # Airflow 启动且插件被加载时执行的回调。 # 注意方法定义中务必带 *args 与 **kwargs以兼容未来 on_load(...) 注入的额外参数。 def on_load(*args, **kwargs): # ... 执行插件启动动作 pass # 可重定向到外部系统的全局 operator extra links以按钮形式出现在任务页。 # 注意全局 extra link 可在单个 operator 层被覆盖。 global_operator_extra_links [] # 覆盖或为既有 Airflow Operator 新增链接的 operator extra links按钮形式。 operator_extra_links [] # 要注册以便在 DAG 中使用的 timetable 类列表 timetables [] # 可作为 DAG 自定义截止时间的 deadline reference 类列表。 # 自定义 deadline 引用类必须在此注册才能在调度器反序列化时被解析 # 未注册的类一旦被 DAG 使用会抛出 DeadlineReferenceNotRegistered。 deadline_references [] # 插件提供的 Listener 列表用于监听特定事件如 TaskInstance 状态变化。Listeners 是 Python 模块。 listeners []你可以通过继承子类化派生它也可把属性定义成 property 以执行额外的初始化。name是必填字段。修改插件后务必重启 webserver 与 scheduler 使其生效。插件管理界面Airflow 3.1 引入了插件管理界面Plugin Management Interface位于 UI 的Admin → Plugins。该页面允许你查看当前已安装的插件是可视化确认插件是否被识别的手段。一个完整的插件示例下面的示例源自 plugins.rst 的完整代码块演示了插件能够注入的全部核心对象# 这是你要派生的插件基类 from airflow.plugins_manager import AirflowPlugin from fastapi import FastAPI from fastapi.middleware.trustedhost import TrustedHostMiddleware # 需要派生的基类 from airflow.hooks.base import BaseHook from airflow.providers.amazon.aws.transfers.gcs_to_s3 import GCSToS3Operator # 通过模板中的 {{ macros.test_plugin.plugin_macro }} 展示 def plugin_macro(): pass # 创建一个集成进 Airflow Rest API 的 FastAPI 应用 app FastAPI() app.get(/) async def root(): return {message: Hello World from FastAPI plugin} app_with_metadata {app: app, url_prefix: /some_prefix, name: Name of the App} # 创建作用于所有 server api 请求的 FastAPI 中间件 middleware_with_metadata { middleware: TrustedHostMiddleware, args: [], kwargs: {allowed_hosts: [example.com, *.example.com]}, name: Name of the Middleware, } # 创建渲染在 Airflow UI 中的外部视图 external_view_with_metadata { # 外部视图名称将显示在 UI 中 name: Name of the External View, # 外部视图的源 URL。URL 可使用上下文变量做模板化——可用变量取决于渲染位置 # 即 (DAG_ID, RUN_ID, TASK_ID, MAP_INDEX, ASSET_ID, ASSET_URI) 的子集 href: https://example.com/{DAG_ID}/{RUN_ID}/{TASK_ID}/{MAP_INDEX}, # 外部视图的挂载位置决定该视图在 UI 中哪里加载。 # 支持 Literal[nav, dag, dag_run, task, task_instance, asset, base]默认 nav destination: dag_run, # 可选图标svg 文件的 url icon: https://example.com/icon.svg, # 可选深色主题图标svg 文件的 url不提供则亮暗主题都使用 icon icon_dark_mode: https://example.com/dark_icon.svg, # 可选参数外部视图渲染的相对 URL 路径。不提供则外部视图渲染为外部链接 # 提供则以 iframe 形式内嵌在 UI 中。不应以斜杠开头 url_route: my_external_view, # 可选类别仅对 destination nav 有效用于把外部链接分组进导航栏。 # 会匹配 [browse, docs, admin, user] 现有菜单未匹配则新建菜单 category: browse, # 可选标记仅对 destination nav 有效。为 True 时该项始终直接渲染在导航工具栏上 # 而不进入 Plugins 子菜单。当存在两个及以上未置顶项时仍会归入子菜单 # 若只剩一个未置顶项也会显示在工具栏。默认 False nav_top_level: True, # 可选作用域限制该视图的显示位置。整体省略则处处显示默认。 # 见下文 Scoping a view to specific Dags and tasks applies_to: { dag_tags: [production, ml], dag_ids: [my_dag, my_other_dag], }, } # 注意React app 集成目前为实验性接口未来可能变化 react_app_with_metadata { # React app 名称将显示在 UI 中 name: Name of the React App, # React app 的 bundle URL即 React app 被服务的地址可为静态文件或 CDN。 # URL 可使用上下文变量做模板化可用变量取决于渲染位置 # 即 (DAG_ID, RUN_ID, TASK_ID, MAP_INDEX, ASSET_ID, ASSET_URI) 的子集 bundle_url: https://example.com/static/js/my_react_app.js, # React app 挂载位置。 # 支持 Literal[nav, dag, dag_run, task, task_instance, asset, base]默认 nav。 # 也可放入既有页面支持视图 [dashboard, dag_overview, task_overview] # 通过 css order 规则决定 flex 顺序来定位元素。 # 使用 base 将 app 挂载到基础布局如工具栏条宿主使用 flex 容器 # 可在根 JSX 中设置 order 控制位置。 destination: task, # 可选图标svg 文件 url icon: https://example.com/icon.svg, # 可选深色主题图标不提供则使用 icon icon_dark_mode: https://example.com/dark_icon.svg, # React app 的 URL 路由相对 Airflow UI 基础 URL不应以斜杠开头 url_route: my_react_app, # 可选类别仅对 destination nav 有效匹配 [browse, docs, admin, user]未匹配则新建菜单 category: browse, # 可选置顶标记仅对 destination nav 有效默认 False nav_top_level: True, # 可选作用域限制 app 显示位置整体省略则处处显示默认 applies_to: { dag_tags: [production, ml], operators: [KubernetesPodOperator], }, } # 定义插件类 class AirflowTestPlugin(AirflowPlugin): name test_plugin macros [plugin_macro] fastapi_apps [app_with_metadata] fastapi_root_middlewares [middleware_with_metadata] external_views [external_view_with_metadata] react_apps [react_app_with_metadata]示例要点解读macrosplugin_macro会以{{ macros.test_plugin.plugin_macro }}形式出现在 DAG 模板中——test_plugin正是插件name在宏命名空间中的命名空间前缀。FastAPI App通过url_prefix挂到/some_prefix路径下成为 Airflow REST API 的一部分。FastAPI 中间件TrustedHostMiddleware会作用于所有 API 请求kwargs中限定只信任example.com及其子域。URL 模板化href/bundle_url中的{DAG_ID}风格占位符会被真实上下文替换占位符可用集取决于视图被渲染的位置。将视图限定到指定 DAG 与任务applies_to 作用域默认情况下external view 或 React app 会出现在其destination匹配的每一个页面上。可选的applies_to块把显示范围收窄使标签页只在相关处出现而不是出现在所有 DAG 上applies_to: { dag_tags: [ml], # DAG 带有任一这些标签 dag_ids: [train_pipeline], # 精确的 dag_id task_ids: [train_model], # 精确的 task_id operators: [KubernetesPodOperator], # operator 类名 }匹配语义与关键规则所有键均可选。operators与operator_names是分开匹配的与任务实例过滤器的语义一致operators匹配 operator 类名而operator_names匹配 UI 中显示的展示名即 operator 的custom_operator_name。对普通 operator 二者相同因此任一键都可以装饰器任务则不同——例如task.bash任务的展示名是task.bash但其私有类名是_BashDecoratedOperator此时要用operator_names才能命中。条件组合方式与 Kubernetes 标签选择器一致——键内 OR键间 ANDDAG 只要命中dag_tags列表中的任一标签即满足该条件而同时配置了dag_tags与operators的视图则要求两者都命中。一个关键设计是AND 只在当前页面能评估的条件之间生效。例如task_ids条件在 DAG 级页面无法判断此时会被跳过而不是导致匹配失败。这样同一个applies_to块可以被插件的 DAG 级与任务级 destination 共用。各 destination 能评估哪些条件官方文档与 源码中的映射表_EVALUABLE_CRITERIA_BY_DESTINATION保持同步Destinationdag_tags/dag_idstask_ids/operators/operator_namesdag、dag_run、dag_overview评估跳过task、task_overview、task_instance评估评估nav、base、dashboard、asset跳过跳过如果某页面上配置的所有条件都无法评估视图照常显示。在 task group 页面上任务级条件会被跳过因为 group 不是 task。非法 applies_to 的处理格式非法的applies_to不是字典、写了未知条件名、或某个条件的值不是字符串列表会在插件加载时被警告并忽略视图仍以无作用域状态加载配置了 destination 无法评估的条件例如在dag视图上写task_ids也会被警告因为它在那里不起作用。相关校验逻辑见 plugins_manager.py 的_validate_applies_to。如果视图作用域不符预期请查看 API server 日志中的这些警告。重要提醒applies_to只是显示层面的便利设施不是授权边界。它只控制 UI 是否展示标签页不控制底层视图是否可被访问——知道url_route的用户仍然可以直接导航到它。要限制谁能查看插件数据请使用访问控制access control。React App 的上下文 propsReact app 集成目前为实验性以下 props 未来可能变化。与外部视图只能在bundle_url中通过{DAG_ID}风格 token 接收上下文不同React app 作为组件被渲染直接以 props 形式接收上下文。可用 props 取决于 app 挂载的位置destination与路由dagId、runId、taskId、mapIndex、assetId—— 当前路由中的标识符字符串存在时提供assetUri—— 挂载在 asset 路由上时当前 asset 的 URIdag、dagRun、taskInstance、asset—— 当前路由对应的完整记录与对应 REST API 响应模式一致DAGDetailsResponse、DAGRunResponse、TaskInstanceResponse、AssetResponse。每个对象只有在路由中出现其依赖的标识符后才提供并且来自详情页已填充的 UI 查询缓存不产生额外请求。在没有这些标识符的路由 /destination上如nav、base、dashboard对应对象为undefined。将视图排除在 CSRF 保护之外我们强烈建议你为所有视图启用 CSRF 保护。但确有需要时可以用装饰器豁免某些视图from airflow.www.app import csrf csrf.exempt def my_handler(): # ... return ok把插件作为 Python 包分发可以通过 setuptools entrypoint 机制加载插件在你的包里用 entrypoint 声明插件。一旦该包被安装Airflow 会自动从 entrypoint 列表加载已注册的插件。注意entrypoint 名称如my_plugin和插件类名都不会影响插件自身的模块名与类名。# my_package/my_plugin.py from airflow.plugins_manager import AirflowPlugin class MyAirflowPlugin(AirflowPlugin): name my_namespace然后在pyproject.toml中声明[project.entry-points.airflow.plugins] my_plugin my_package.my_plugin:MyAirflowPlugin这是把插件打包、安装、随 Python 环境自动发现的最规范方式适合插件要跨多个 Airflow 部署复用的场景。Airflow 3 中的 Flask AppBuilder 与 Flask Blueprint 变迁Airflow 2 的插件支持 Flask AppBuilder 视图appbuilder_views、Flask AppBuilder 菜单项appbuilder_menu_items以及 Flask Blueprintflask_blueprints。在 Airflow 3 中这些已被新的External Viewsexternal_views、FastAPI Appsfastapi_apps、FastAPI Middlewaresfastapi_root_middlewares和 React Appsreact_apps取代新接口提供了更强的功能与更好的 Airflow UI 集成。所有新插件都应使用新接口。不过为平滑迁移到 Airflow 3社区提供了 Flask / FAB 插件的兼容层只需安装 FAB provider并按照 Airflow 3 迁移指南调整代码即可继续使用现有的 Flask AppBuilder 视图、Flask Blueprint 和 Flask AppBuilder 菜单项。插件故障排查除了前面提到的airflow plugins命令你还可以借助 Flask CLI 排查问题。运行前需要设置环境变量FLASK_APP指向airflow.www.app:create_app。例如打印所有路由FLASK_APPairflow.www.app:create_app flask routes这会输出应用注册的全部路由有助于确认插件注册的 FastAPI 端点 / 视图路由是否真正进入了 Flask 应用的路由表。小结与最佳实践结合本仓库源码与官方文档编写与维护插件时可遵循以下清单命名唯一name必填且全局唯一否则会在 加载阶段 被以「already registered」跳过写新代码用新接口Airflow 3 中一律使用external_views/fastapi_apps/fastapi_root_middlewares/react_appsFlask/FAB 接口仅供既有插件迁移过渡on_load带上*args, **kwargs以兼容未来注入的新参数视图有副作用就用applies_to收窄范围但不要把它当作安全边界敏感数据仍需访问控制记得重启进程修改插件后重启 webserver / scheduler任务内使用的插件在 fork 模式下要重启 worker / scheduler 才更新或开启core.execute_tasks_new_python_interpreter True多部署复用时打包成 entrypoint 插件通过[project.entry-points.airflow.plugins]声明安装即发现排障三件套airflow plugins查看加载信息、Admin → Plugins查看 UI 插件管理界面、FLASK_APPairflow.www.app:create_app flask routes检查路由并留意 API server 日志中与applies_to相关的警告。在此基础上你便可以用纯「放置文件 / 注册 entrypoint」的方式为 Airflow 打造 Hive 元数据浏览器、异常检测、审计、SLA 监控等贴合自身生态的扩展而无需 fork 或改动核心代码。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价