news 2026/9/10 23:33:39

Apache Airflow 新增 REST API 端点完全指南:从 FastAPI 路由实现到 prek 钩子验证

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Apache Airflow 新增 REST API 端点完全指南:从 FastAPI 路由实现到 prek 钩子验证

Apache Airflow 新增 REST API 端点完全指南:从 FastAPI 路由实现到 prek 钩子验证

【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow

本指南面向希望在 Apache Airflow 3(当前仓库)的 REST API 中新增端点的开发者,完整讲解从接口设计决策、FastAPI 路由实现、Pydantic 数据模型定义、单元测试编写,到 prek 钩子自动更新 OpenAPI 规范的端到端流程。读完本文,你将掌握publicui两套 API 的选型标准,能够基于 AirflowRouter 注册带权限校验与查询参数的路由,并让新端点自动进入 v2-rest-api-generated.yaml 与交互式文档。

一、先理解 Airflow 3 的 API 分层:public 与 ui

Airflow 3 的 REST API 全部构建在 FastAPI 之上,核心源码位于airflow-core/src/airflow/api_fastapi/core_api,其中:

  • routes/public:公共 API 端点,路径以/api/v2为前缀。它们经过标准化、文档完善,且保证向后兼容。外部用户、SDK 与集成方应只依赖这一层。
  • routes/ui:专门为前端(Airflow UI)定制的端点,路径以/ui为前缀,不承诺向后兼容,可随时根据前端需要调整,外部不应依赖。

该分层在仓库中有着明确的落地证据:生成的 OpenAPI 规范文件 v2-rest-api-generated.yaml 的info.description中写道:/api/v2下的端点"可以安全使用、稳定且向后兼容",而/ui下的端点"专为 UI 服务,可能随前端需求发生破坏性变更"。同时,在 ui/init.py 中,ui_router注册时显式设置了include_in_schema=False,这意味着 UI 端点不会出现在公开的 OpenAPI schema 文档中。

选型建议:新增端点时,应尽最大可能使其可被社区复用,设计稳定、标准化,因此优先放入public;只有当数据类型过于特殊或频繁变动(典型如 Grid、Gantt、Calendar 等页面专用数据结构)时,才放入ui。仓库中的 public 路由目录(含dags.pyconnections.pyvariables.pytask_instances.pyxcom.py等 31 个模块)与 ui 路由目录(含grid.pygantt.pycalendar.pydashboard.pydependencies.py等 16 个模块)的对比,就是这一决策的直观样例。

二、Step 1:实现端点逻辑

2.1 确定接口归属并创建路由

根据上面的分层决策,导航到airflow-core/src/airflow/api_fastapi/core_api/routes/publicairflow-core/src/airflow/api_fastapi/core_api/routes/ui下的对应模块,使用AirflowRouter注册新路由。AirflowRouter是 FastAPIAPIRouter的轻量扩展,位于 airflow-core/src/airflow/api_fastapi/common/router.py,其核心差异在于:若不显式传operation_id,会自动使用函数名作为operation_id,这保证了 OpenAPI 规范中操作标识符的确定性。

原文档给出的路由骨架如下:

@dags_router.get("/dags") # permissions go in the dependencies parameter here async def get_dags( *, limit: int = 100, offset: int = 0, tags: Annotated[list[str] | None, Query()] = None, dag_id_pattern: str | None = None, only_active: bool = True, paused: bool | None = None, order_by: str = "dag_id", session: SessionDep, ) -> DagCollectionResponse: pass

2.2 对照真实源码理解各要素

仓库中 public/dags.py 的get_dags是该骨架在生产环境的完整形态,可逐项对照:

dags_router = AirflowRouter(tags=["DAG"], prefix="/dags") @dags_router.get("", dependencies=[Depends(requires_access_dag(method="GET"))]) def get_dags( limit: QueryLimit, offset: QueryOffset, tags: QueryTagsFilter, dag_id_pattern: QueryDagIdPatternSearch, paused: QueryPausedFilter, ... order_by: Annotated[SortParam, Depends(...)], readable_dags_filter: ReadableDagsFilterDep, session: SessionDep, ) -> DAGCollectionResponse: ...

要点拆解:

  • HTTP 方法与路径@dags_router.get("")配合prefix="/dags",实际对外路径为/api/v2/dags。同文件中还演示了getpatchpostdelete各种方法的使用(如patch_dag@dags_router.patch("/{dag_id}")favorite_dag@dags_router.post("/{dag_id}/favorite"))。
  • 权限声明:放在装饰器的dependencies参数中。requires_access_dag(method="GET")是 security.py 提供的工厂函数,它从请求的 path/query 中提取dag_id,调用 AuthManager 的is_authorized_dag完成鉴权;写操作需对应method="PUT"method="DELETE"
  • 查询参数的类型化:生产代码不使用裸类型,而是复用airflow.api_fastapi.common.parameters中的预定义类型,例如QueryLimitQueryOffsetQueryTagsFilterQueryDagIdPatternSearch,它们封装了默认值、校验规则与 OpenAPI 描述;排序参数通过SortParam(...).dynamic_depends()动态构造Depends
  • 可读性过滤ReadableDagsFilterDep(对应 security.py 的ReadableDagsFilterDep = Annotated[PermittedDagFilter, Depends(permitted_dag_filter_factory("GET"))])会将当前用户无权读取的 DAG 在 SQL 层面过滤掉,保证"存在性不泄露"。
  • 数据库会话session: SessionDep是 FastAPI 依赖注入的 SQLAlchemy 会话,所有查询都经由它执行,配合paginated_select统一实现分页与total_entries统计。
  • 返回类型注解-> DagCollectionResponse是必须的,FastAPI 依据该注解生成 OpenAPI 的 response schema。

2.3 错误响应的 OpenAPI 文档化

对于可能抛出多种 HTTP 异常的端点,仓库推荐使用create_openapi_http_exception_doc显式声明响应码,例如同文件中get_dag的定义:

@dags_router.get( "/{dag_id}", responses=create_openapi_http_exception_doc( [status.HTTP_400_BAD_REQUEST, status.HTTP_404_NOT_FOUND, HTTP_422_UNPROCESSABLE_CONTENT] ), dependencies=[Depends(requires_access_dag(method="GET"))], )

该工具函数位于 core_api/openapi/exceptions.py,它会为每个列出的状态码生成符合 Airflow 统一错误结构的文档,让 API 消费者提前知晓错误语义(如 404 表示资源不存在、409 表示存在冲突如"任务实例仍在运行")。

三、Step 2:为端点编写测试

3.1 手动验证

在编写测试前,先在本地启动 API 服务,用 curl 或 FastAPI 自带文档手动验证端点行为符合预期,包括查询参数组合、权限拦截与异常分支。

3.2 初始化单元测试

测试目录结构与源码一一对应:

  • 公共端点测试位于 airflow-core/tests/unit/api_fastapi/core_api/routes/public,每个端点模块都有对应测试文件,如test_dags.pytest_connections.pytest_variables.pytest_xcom.pytest_task_instances.py等 30 个文件;
  • UI 端点测试位于 airflow-core/tests/unit/api_fastapi/core_api/routes/ui。

测试应覆盖以下维度(原文档要求"extensive tests"):

  • 查询参数:合法值、边界值、非法值(应返回 422 校验错误);
  • 权限:无认证访问返回 401,无授权访问返回 403,不同 method 的权限差异;
  • 错误处理:资源不存在返回 404、并发/状态冲突返回 409 等。

3.3 全局路由约束测试

仓库还有一个值得借鉴的全局测试 test_routes.py,它从public_routerauthenticated_router的注册表出发做结构性断言:

  • test_no_auth_routes:验证除NO_AUTH_PATHS/api/v2/auth/login/api/v2/auth/logout/api/v2/version/api/v2/monitor/health)外,所有路由都必须经过认证;
  • test_routes_with_responses:验证每个受保护路由都声明了 401/403 响应,防止新增路由时意外漏掉鉴权。

这意味着:你新增的受保护端点若未配置 401/403 响应文档,这一套测试会直接失败——这是 Airflow 通过测试强制保障 API 安全一致性的机制。

四、Step 3:文档自动生成与核对

Airflow 的 API 文档由 FastAPI 自动生成,并经过 prek 钩子(见下一节)固化。新增端点后:

  1. 启动服务并访问/docs(FastAPI 内置 Swagger UI);
  2. 核对新端点的请求体类型、返回类型、查询参数、校验规则是否清晰、完整、符合预期;
  3. 确认错误响应码(400/401/403/404/409/422)已在文档中正确呈现。

在仓库中,这份自动生成的交互式文档的数据源就是 v2-rest-api-generated.yaml(当前仓库中约 1.7 万行、覆盖所有 public 端点),它同时被前端代码生成工具消费。

五、Step 4:运行 prek 钩子完成质量门禁

5.1 prek 是什么

prek是 Apache Airflow 项目的静态检查与代码生成工具链(其依赖声明可见于 dev/breeze/pyproject.toml 中的"prek>=0.4.14"),涵盖静态代码检查、格式化以及各类生成文件的校验。新增 API 端点后必须运行它以通过 CI 门禁。

执行全部 prek 钩子:

prek --all-files

5.2 OpenAPI 规范的持久化更新

新增端点后,持久化的 OpenAPI 规范文件v2-rest-api-generated.yaml需要同步更新。这一步由专用的 prek 钩子自动完成,你只需:

  1. 运行 prek 钩子,它会基于当前 FastAPI 应用重新生成规范文件;
  2. 将生成的变更git add并提交。

该钩子的实现链路是:scripts/ci/prek/generate_openapi_spec.py 通过 Breeze 容器调用 scripts/in_container/run_generate_openapi_spec.py,后者以SimpleAuthManager初始化 FastAPI 应用并调用generate_openapi_file(app=create_app(), file_path=OPENAPI_SPEC_FILE)写回 v2-rest-api-generated.yaml,同时生成用于 UI 代码生成的私有规范_private_ui.yaml与 SimpleAuthManager 专属规范。这意味着你无需手写任何 OpenAPI YAML,规范文件始终与代码保持同步。

六、可选:添加 Pydantic 模型

当新端点涉及全新的数据结构时,需要定义对应的 Pydantic 模型,用于校验、序列化/反序列化请求与响应。原文档给出的最小示例:

class DagModelResponse(BaseModel): """Dag serializer for responses.""" dag_id: str dag_display_name: str is_paused: bool is_active: bool last_parsed_time: datetime | None

6.1 生产环境中的模型形态

仓库中的模型比示例更复杂、更具参考性。响应模型统一放在airflow-core/src/airflow/api_fastapi/core_api/datamodels下,以 datamodels/dags.py 为例:

  • DAGResponse:核心序列化模型,声明dag_iddag_display_nameis_pausedis_stalelast_parsed_time等 20+ 字段;通过@field_validatorowners由逗号分隔字符串归一为列表、@computed_field计算is_backfillablefile_token;还通过AliasGenerator处理响应字段与模型字段的命名映射(如next_dagrun_logical_datenext_dagrun)。
  • DAGCollectionResponse:集合响应统一为dags: Iterable[DAGResponse]+total_entries: int的结构,这是全仓库所有列表端点(Connections、Variables、Assets 等)的通用约定,可从 datamodels 目录 中大量XxxCollectionResponse得到印证。
  • DAGPatchBody/DAGPatchBodyPartial:请求体模型,DAGPatchBodyPartialmake_partial_model工具生成"全部字段可选"的变体,用于支持update_mask局部更新语义。

6.2 模型自动进入 OpenAPI

这些模型只要被某个端点实际引用(作为参数或返回类型注解),就会自动出现在 OpenAPI 规范文件中,无需手动登记。因此:

  1. 新增或修改模型后,重新运行 prek 钩子以更新所有生成文件(包括v2-rest-api-generated.yaml中的components/schemas部分);
  2. 未被任何端点引用的模型不会进入规范——这保证了规范的整洁性。

七、把新端点接入应用:路由注册机制

原文档没有展开、但对实际提交 PR 至关重要的环节是:仅定义路由函数还不够,必须把 Router 挂载到应用上。完整链路如下:

  1. routes/public/your_module.py中创建xxx_router = AirflowRouter(...)
  2. 在 routes/public/init.py 中导入该 router,并通过authenticated_router.include_router(xxx_router)挂载——authenticated_router在路由器级别注入了Depends(get_user),作为"防呆兜底",即使某个路由忘了写自己的鉴权依赖,也会被强制要求认证(该设计动机在文件注释中有明确说明);
  3. public_router = AirflowRouter(prefix="/api/v2")再统一挂载authenticated_router,以及无需认证的monitor_routerversion_routerauth_router
  4. 最终由 core_api/app.py 的init_views调用app.include_router(ui_router)app.include_router(public_router)完成应用级注册,同时它还注册了/api/v1与未知/api/*的 404 兜底路由(明确提示/api/v1在 Airflow 3 中已被移除、应改用/api/v2)。

同理,UI 端点的注册路径是 routes/ui/init.py:ui_router = AirflowRouter(prefix="/ui", include_in_schema=False, dependencies=[Depends(get_user)]),同样以 router 级依赖保证全部 UI 路由必须认证。

八、安全基线:新端点必须考虑的三件事

结合 core_api/security.py 与 test_routes.py,新增端点时请默念以下安全清单:

  1. 认证兜底:通过authenticated_router(public)或 router 级Depends(get_user)(ui)挂载,或至少声明Depends(get_user),使未认证请求返回 401;
  2. 细粒度授权:DAG 类资源使用requires_access_dag(method=...)并按读写方法区分;非 DAG 资源(Pool、Connection、Variable 等)使用各自的requires_access_*工厂;列表类端点注入Readable*FilterDep做行级过滤;
  3. 响应码声明:用create_openapi_http_exception_doc声明 401/403 及业务错误码,否则test_routes_with_responses这类全局测试会拦截合入。

九、小结:新增端点的完整检查清单

按原文档流程并结合仓库实践,一个 API 端点的合入路径可归纳为:

  1. 判断归属:标准化、向后兼容 →public;UI 专属、易变 →ui
  2. 在对应routes/目录实现路由函数:选好 HTTP 方法、查询参数类型、dependencies权限、responses错误码、Pydantic 返回类型注解;
  3. 将 Router 挂载到 public/init.py 或 ui/init.py;
  4. 在对应tests/unit/api_fastapi/core_api/routes/下补齐测试:查询参数、权限、错误处理全覆盖;
  5. 访问/docs核对文档(body/return 类型、query 参数、校验规则);
  6. 运行prek --all-files,让钩子自动更新 v2-rest-api-generated.yaml 等生成文件,git add后提交。

如果你在此基础上进一步调整了 Airflow 的整体架构,可参考仓库中的 架构图绘制指南 来同步更新架构文档。

【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/10 23:29:56

武汉青山区洗衣机维修推荐,欧米到家解决老旧洗衣机维修和保养需求

前言洗衣机是现代家庭使用频率较高的家电之一,长期运行后容易出现不脱水、不排水、不进水、漏水、异响、无法启动、显示故障代码等问题。尤其武汉地区家庭使用洗衣机频率较高,面对设备突然故障时,选择专业、规范的维修服务非常重要。欧米到家…

作者头像 李华
网站建设 2026/9/10 23:26:44

如何安装 TVBoxOSC:电视盒子管理工具的部署与配置指南

如何安装 TVBoxOSC:电视盒子管理工具的部署与配置指南 【免费下载链接】TVBoxOSC TVBoxOSC - 一个基于第三方项目的代码库,用于电视盒子的控制和管理。 项目地址: https://gitcode.com/GitHub_Trending/tv/TVBoxOSC TVBoxOSC 是一个面向电视盒子控…

作者头像 李华
网站建设 2026/9/10 23:25:42

关于墨衍 MoGrow 的 15 个常见问题(FAQ 合集)

墨衍 MoGrow 是面向开发者与企业的一站式 AI 数字营销平台,提供 AI 选题创作、多平台一键分发与 SEO & GEO 双端优化。本文把用户最常问的 15 个问题集中回答,便于一次性建立完整认知。 一、产品认知类 1. 墨衍 MoGrow 是什么? 墨衍 MoGr…

作者头像 李华
网站建设 2026/9/10 23:23:43

SpringBoot电竞商城系统架构与高并发实践

1. 项目概述:电竞周边商城的商业与技术价值这个基于SpringBoot的游戏周边商城系统,本质上是一个垂直领域的电商平台,专门服务于快速增长的电竞衍生品市场。根据最新行业报告,全球电竞周边市场规模已突破50亿美元,年增长…

作者头像 李华