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 规范的端到端流程。读完本文,你将掌握public与ui两套 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.py、connections.py、variables.py、task_instances.py、xcom.py等 31 个模块)与 ui 路由目录(含grid.py、gantt.py、calendar.py、dashboard.py、dependencies.py等 16 个模块)的对比,就是这一决策的直观样例。
二、Step 1:实现端点逻辑
2.1 确定接口归属并创建路由
根据上面的分层决策,导航到airflow-core/src/airflow/api_fastapi/core_api/routes/public或airflow-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: pass2.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。同文件中还演示了get、patch、post、delete各种方法的使用(如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中的预定义类型,例如QueryLimit、QueryOffset、QueryTagsFilter、QueryDagIdPatternSearch,它们封装了默认值、校验规则与 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.py、test_connections.py、test_variables.py、test_xcom.py、test_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_router与authenticated_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 钩子(见下一节)固化。新增端点后:
- 启动服务并访问
/docs(FastAPI 内置 Swagger UI); - 核对新端点的请求体类型、返回类型、查询参数、校验规则是否清晰、完整、符合预期;
- 确认错误响应码(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-files5.2 OpenAPI 规范的持久化更新
新增端点后,持久化的 OpenAPI 规范文件v2-rest-api-generated.yaml需要同步更新。这一步由专用的 prek 钩子自动完成,你只需:
- 运行 prek 钩子,它会基于当前 FastAPI 应用重新生成规范文件;
- 将生成的变更
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 | None6.1 生产环境中的模型形态
仓库中的模型比示例更复杂、更具参考性。响应模型统一放在airflow-core/src/airflow/api_fastapi/core_api/datamodels下,以 datamodels/dags.py 为例:
DAGResponse:核心序列化模型,声明dag_id、dag_display_name、is_paused、is_stale、last_parsed_time等 20+ 字段;通过@field_validator将owners由逗号分隔字符串归一为列表、@computed_field计算is_backfillable与file_token;还通过AliasGenerator处理响应字段与模型字段的命名映射(如next_dagrun_logical_date→next_dagrun)。DAGCollectionResponse:集合响应统一为dags: Iterable[DAGResponse]+total_entries: int的结构,这是全仓库所有列表端点(Connections、Variables、Assets 等)的通用约定,可从 datamodels 目录 中大量XxxCollectionResponse得到印证。DAGPatchBody/DAGPatchBodyPartial:请求体模型,DAGPatchBodyPartial由make_partial_model工具生成"全部字段可选"的变体,用于支持update_mask局部更新语义。
6.2 模型自动进入 OpenAPI
这些模型只要被某个端点实际引用(作为参数或返回类型注解),就会自动出现在 OpenAPI 规范文件中,无需手动登记。因此:
- 新增或修改模型后,重新运行 prek 钩子以更新所有生成文件(包括
v2-rest-api-generated.yaml中的components/schemas部分); - 未被任何端点引用的模型不会进入规范——这保证了规范的整洁性。
七、把新端点接入应用:路由注册机制
原文档没有展开、但对实际提交 PR 至关重要的环节是:仅定义路由函数还不够,必须把 Router 挂载到应用上。完整链路如下:
- 在
routes/public/your_module.py中创建xxx_router = AirflowRouter(...); - 在 routes/public/init.py 中导入该 router,并通过
authenticated_router.include_router(xxx_router)挂载——authenticated_router在路由器级别注入了Depends(get_user),作为"防呆兜底",即使某个路由忘了写自己的鉴权依赖,也会被强制要求认证(该设计动机在文件注释中有明确说明); public_router = AirflowRouter(prefix="/api/v2")再统一挂载authenticated_router,以及无需认证的monitor_router、version_router、auth_router;- 最终由 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,新增端点时请默念以下安全清单:
- 认证兜底:通过
authenticated_router(public)或 router 级Depends(get_user)(ui)挂载,或至少声明Depends(get_user),使未认证请求返回 401; - 细粒度授权:DAG 类资源使用
requires_access_dag(method=...)并按读写方法区分;非 DAG 资源(Pool、Connection、Variable 等)使用各自的requires_access_*工厂;列表类端点注入Readable*FilterDep做行级过滤; - 响应码声明:用
create_openapi_http_exception_doc声明 401/403 及业务错误码,否则test_routes_with_responses这类全局测试会拦截合入。
九、小结:新增端点的完整检查清单
按原文档流程并结合仓库实践,一个 API 端点的合入路径可归纳为:
- 判断归属:标准化、向后兼容 →
public;UI 专属、易变 →ui; - 在对应
routes/目录实现路由函数:选好 HTTP 方法、查询参数类型、dependencies权限、responses错误码、Pydantic 返回类型注解; - 将 Router 挂载到 public/init.py 或 ui/init.py;
- 在对应
tests/unit/api_fastapi/core_api/routes/下补齐测试:查询参数、权限、错误处理全覆盖; - 访问
/docs核对文档(body/return 类型、query 参数、校验规则); - 运行
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),仅供参考