在Django REST API视图中调用时动态创建Airflow DAG的问题咨询
动态生成Airflow DAG的实现方案
原有方案失效原因
- Airflow的DAG是由Scheduler进程定期扫描
dags_folder目录下的Python文件、解析代码生成的,你在Django视图中直接实例化DAG只会存储在Django进程的内存中,Airflow的Scheduler完全无法感知,自然不会生效。 - 将
dags_folder修改为Django views.py路径的操作也无效:Scheduler扫描时只会执行views.py的顶层代码,不会触发API请求对应的业务逻辑,不可能生成你需要的动态DAG。
两种主流实现方案
方案1:Django接口直接生成DAG文件到dags目录
这是最通用的实现方式,接口拿到作业元数据后,直接生成标准的DAG代码写入Airflow的dags_folder目录,Scheduler下次扫描就会自动加载该DAG。
示例代码如下:
import os from django.conf import settings from rest_framework.response import Response def create_job_view(request): # 此处省略你的业务逻辑:获取前端参数、写入数据库等 dag_id = "test_dag_1" schedule = "@daily" owner = "airflow" start_date = "2021, 9, 13" # 构造DAG代码字符串 dag_code = f"""from airflow import DAG from airflow.operators.python import PythonOperator import datetime def hello_world_py(): print("hello world") default_args = {{ 'owner': '{owner}', 'start_date': datetime.datetime({start_date}) }} with DAG( '{dag_id}', schedule_interval='{schedule}', default_args=default_args, catchup=False ) as dag: t1 = PythonOperator( task_id='hello_world', python_callable=hello_world_py ) globals()['{dag_id}'] = dag """ # 写入到Airflow的dags目录 dag_file_path = os.path.join(settings.AIRFLOW_DAGS_FOLDER, f"{dag_id}.py") with open(dag_file_path, "w", encoding="utf-8") as f: f.write(dag_code) return Response({"status": "success", "dag_id": dag_id})
注意事项:
- 提前在Django配置中定义
AIRFLOW_DAGS_FOLDER为你的Airflow dags目录绝对路径 - 给Django运行进程开放dags目录的读写权限
- 可调整Airflow配置项
dag_dir_list_interval降低扫描间隔,让新DAG更快被识别 - 生成代码时注意f-string转义,字典的大括号需要双写避免解析错误
- 如需更新/删除DAG,直接修改/删除dags目录下对应的
.py文件即可
方案2:单文件动态生成所有DAG
如果所有DAG的结构高度相似,只是参数不同,可以提前在dags目录下写一个通用DAG生成脚本,直接从Django的数据库读取作业元数据循环生成DAG,不需要Django写文件,接口只需要把元数据存入数据库即可。
示例代码(dags目录下的dynamic_dags.py):
from airflow import DAG from airflow.operators.python import PythonOperator import datetime import os import sys # 加载Django环境 sys.path.append("/path/to/your/django/project/root") os.environ.setdefault("DJANGO_SETTINGS_MODULE", "your_project.settings") import django django.setup() from yourapp.models import JobMetadata def hello_world_py(): print("hello world") # 查询所有启用的作业元数据 active_jobs = JobMetadata.objects.filter(is_active=True) for job in active_jobs: dag_id = job.dag_id default_args = { "owner": job.owner, "start_date": job.start_date } # 生成DAG with DAG( dag_id, schedule_interval=job.schedule, default_args=default_args, catchup=False ) as dag: t1 = PythonOperator( task_id="hello_world", python_callable=hello_world_py ) # 加入全局变量让Airflow识别 globals()[dag_id] = dag
注意事项:
- 保证Airflow运行环境安装了Django项目的所有依赖,且有权限访问Django的数据库
- 作业数量较多时会增加Scheduler的扫描耗时,适合作业量不大的场景
内容的提问来源于stack exchange,提问作者Sony Khan
相关产品推荐
相关产品推荐

