airflow-rest-api-plugin
把 Airflow CLI 命令变成 REST API,让任意应用都能操控 DAG 任务调度
加载项目详情…
本应用为开源项目,仅供学习研究,请遵守其开源协议。
把 Airflow CLI 命令变成 REST API,让任意应用都能操控 DAG 任务调度
加载项目详情…
本应用为开源项目,仅供学习研究,请遵守其开源协议。
凌晨两点,数据工程师小张发现上游数据源异常,需要紧急暂停 ETL DAG 并手动触发修复流程。他打开终端,熟练地敲下 airflow pause 命令,却发现这个操作只能在有 SSH 权限的服务器上执行——而他此刻正在回家的地铁上。
这正是 airflow-rest-api-plugin 要解决的问题:把 Airflow 的命令行能力,以 REST API 的形式暴露出来,让任何能发起 HTTP 请求的应用都能操控 DAG 和任务。它就像给 Airflow 安装了一个通用的遥控面板,从此调度系统的控制权不再被终端窗口所束缚。
Apache Airflow 诞生于 Airbnb,是目前最流行的开源工作流调度平台。其核心调度逻辑通过 Python CLI 实现,Web UI 提供可视化监控,但在 API 化方面长期存在短板。Airflow 1.x 时代,官方并未提供开箱即用的 REST 接口,企业要实现程序化调度通常面临两条路:直接操作底层数据库,或自行封装 CLI 调用。
teamclairvoyant(一家专注大数据和调度系统的技术公司)正是看到了这个痛点,于 2017 年左右推出了这款插件,将 Airflow CLI 的全部能力通过 RESTful 端点一一映射,彻底解决了"命令行调度"的最后一公里问题。
该插件提供了 30+ 个 REST 端点,完整覆盖了 Airflow CLI 的核心命令。按功能可划分为以下几类:
DAG 管控:暂停(pause)、恢复(unpause)、触发(trigger_dag)、查看状态(dag_state)。这些是最常用的运维操作,通过 API 可以实现外部告警系统直接暂停 DAG,或在数据管道断裂时由上游应用自动触发重跑。
任务操作:查询任务失败依赖(task_failed_deps)、测试单个任务(test)、标记成功(mark_success)。开发者可以在任务失败时通过 API 自动执行补救逻辑,无需人工介入。
资源管理:变量(variables)的增删改查、连接(connections)的管理、池(pool)的配置。这些 API 使得配置管理工作可以完全程序化,适合 DevOps 场景下的自动化配置分发。
系统信息:version 查询、list_dags、list_tasks、serve_logs 等。外部监控系统可以通过这些 API 获取 Airflow 运行状态,对接 Prometheus/Grafana 等监控体系。
部署管理:backfill(回填历史)、deploy_dag(热加载 DAG)、refresh_all_dags(刷新 DAG 定义)。这些是高级运维能力,支持不停服更新 DAG 逻辑。
所有 API 均返回统一结构的 JSON 响应,包含状态码、执行时间、输出内容和原始参数,便于客户端做统一的错误处理和日志追踪。
插件代码约 62KB,核心逻辑集中在 plugins/rest_api_plugin.py,采用三层架构设计:
AirflowPlugins 接口层
└── REST_API_Plugin (注册到 Airflow 插件系统)
├── REST_API (Flask-Admin 视图, 提供 Web UI)
├── Flask Blueprint (提供 REST API)
└── REST_API_Response_Util (统一响应格式)
REST_API_Plugin 是核心入口,继承 AirflowPlugin,通过 appbuilder_views、admin_views 和 flask_blueprints 三个渠道注册到 Airflow。Flask Blueprint 是 REST API 的真正载体,接收 HTTP 请求并转发给内部的 process_request 方法处理。
REST_API 类继承自 get_baseview()(Flask-Admin 的 BaseView),负责渲染内置的 Admin 管理页面。每个 API 端点通过解析 apis_metadata 列表动态生成——这个列表中每个元素描述一个 API 的元数据(名称、描述、HTTP 方法、参数),处理函数则通过统一的反射机制执行对应的 CLI 命令。
具体执行流程是:process_request 接收请求后,从 apis_metadata 中找到匹配的 API 定义,取出 cli_program_name 和 cli_command_name,构造出完整的 CLI 命令字符串,然后通过 Python 的 subprocess 模块调用 Airflow CLI,返回结果经格式化后以 JSON 形式返回给客户端。
这种"元编程 + 反射"的架构是整个插件最精妙的设计:用一份 API 元数据同时驱动 Web UI 和 REST API,两套界面共用同一套业务逻辑,后续新增 API 只需在 apis_metadata 列表中添加一项,无需写重复的视图代码。
插件支持两套认证机制:
HTTP Token 认证(简单模式):在 airflow.cfg 中配置 rest_api_plugin_http_token_header_name 和 rest_api_plugin_expected_http_token,客户端请求时在 HTTP Header 中携带 token。这是轻量级的预共享密钥方案,适合内网环境。
JWT 认证(RBAC 模式):当 Airflow 启用 RBAC 时(rbac = True),插件利用 Flask-JWT-Extended 提供标准的 JWT 访问令牌认证。客户端先调用 /api/v1/security/login 获取 access_token 和 refresh_token,后续请求在 Authorization: Bearer <token> Header 中携带 access_token。JWT 默认有效期 15 分钟,支持 refresh_token 续期,最长 30 天。
两种认证方案可独立启用或同时启用,插件会自动检测 Airflow 的 RBAC 配置状态并选择合适的认证方式。
最佳应用场景:构建事件驱动的数据管道(外部数据到达时触发 DAG)、实现 CI/CD 化的 DAG 部署流程、将 Airflow 调度集成到企业内部 IT 服务台系统、对接告警平台的自动恢复机制。
需要正视的局限:该插件仅支持 Airflow 1.x,对 Airflow 2.x 的支持需要另寻方案(Airflow 2.x 官方已内置 REST API)。代码中使用 subprocess 调用 CLI 命令而非直接调用 Python API,存在命令注入的理论风险——虽然 token 认证可以缓解,但高安全要求场景需额外审计。另外,插件未提供 OpenAPI/Swagger 文档,集成方需要自行查阅 README 或抓包分析接口规范。
该插件代表了 Airflow 1.x 时代社区对"调度能力开放"的核心探索方向。它的元数据驱动 API 设计思路影响了许多后来者,即使在 Airflow 2.x 官方推出 REST API 之后,该插件仍在大量运行 Airflow 1.x 的遗留系统中发挥作用。GitHub 325 颗星、92 个 fork 的数据说明它解决了一个真实的工程痛点——不是"要不要开放 API",而是"如何在不改动核心代码的前提下快速开放 API"。这种"插件化扩展"的思路,对于构建可插拔的企业级系统具有普遍参考价值。