如何通过Fast API接口向Docker部署的Apache Airflow提交新DAG
FastAPI 提交 DAG 到 Airflow 可行方案说明
核心结论
调用Airflow CLI的方式可以实现,除此之外还有两种更适配不同场景的方案,生产环境更推荐用官方REST API,测试环境可以选更轻量化的方案。
方案1:调用Airflow CLI实现
前提是两个Docker容器部署在同一宿主机,或者你有跨节点的容器命令执行权限:
- 给FastAPI容器挂载宿主机的
docker.sock,在FastAPI业务逻辑中执行docker exec <your-airflow-scheduler-container-name> airflow dags import <target-dag-path>命令即可完成提交 - 也可以提前把要提交的DAG文件传输到Airflow容器的
dags/目录下,再用CLI命令触发加载
注意:该方案逻辑简单无需额外适配,但开放docker.sock权限有较高的安全风险,仅建议内部可信环境使用,同时要注意CLI版本和部署的Airflow版本严格匹配,避免命令不兼容。
方案2:调用Airflow官方REST API(生产环境首推)
Airflow 2.0+已经内置稳定的DAG管理REST接口,不需要给FastAPI开放过高的容器权限,安全可控:
- 第一步先在Airflow配置中开启REST API,创建拥有DAG编辑权限的服务账号,生成认证Token或直接用账号密码鉴权
- FastAPI侧接收用户提交的DAG文件后,调用Airflow的
POST /api/v1/dags/<dag_id>接口,把DAG文件作为二进制内容提交即可,如需更新已有DAG可调用PATCH接口 - 可以额外调用
POST /api/v1/dags/<dag_id>/dagRuns接口,实现提交DAG后直接触发首次运行的需求
示例FastAPI侧调用代码:
import requests from fastapi import FastAPI, UploadFile app = FastAPI() AIRFLOW_API_URL = "http://你的Airflow-webserver服务地址:8080/api/v1" # 生产环境请把账号密码放到环境变量或配置中心,不要硬编码 AIRFLOW_AUTH = ("admin", "你的Airflow密码") @app.post("/submit-dag") async def submit_new_dag(dag_file: UploadFile): dag_content = await dag_file.read() dag_id = dag_file.filename.rsplit(".", 1)[0] response = requests.post( url=f"{AIRFLOW_API_URL}/dags/{dag_id}", auth=AIRFLOW_AUTH, files={"file": (dag_file.filename, dag_content, "text/python")} ) response.raise_for_status() return {"status": "success", "dag_id": dag_id, "airflow_response": response.json()}
注意:提前配置Airflow的CORS规则允许FastAPI服务地址访问,做好接口鉴权信息的安全存储。
方案3:共享DAG目录挂载(轻量化测试环境首选)
如果两个容器部署在同一宿主机或有共享存储,这个方案实现成本最低:
- 把宿主机的同一个目录,分别挂载到FastAPI容器的上传目录,和Airflow所有服务(scheduler、webserver、worker)的
/opt/airflow/dags/目录 - FastAPI只需要把用户提交的DAG文件写入到这个共享目录即可,Airflow的scheduler默认会每隔30秒(可通过
dag_dir_list_interval配置修改)扫描一次DAG目录,自动加载新的DAG文件
注意:提交前做好DAG文件的语法校验,避免错误的DAG文件导致Airflow加载器报错,多节点Airflow集群需要用NFS、对象存储等共享存储来挂载统一的DAG目录。
选型建议
- 内部测试、快速验证需求:选共享目录挂载方案,实现成本最低
- 生产环境、权限管控要求高:选官方REST API方案,安全稳定易扩展
- 已有成熟容器运维体系、做了Docker sock权限隔离:可选CLI调用方案
内容的提问来源于stack exchange,提问作者ppostnov
相关产品推荐
相关产品推荐

