You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.04 13:42:02