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

使用Jinja生成Airflow模板化DAG:如何保留dag_run.conf等内置字段

问题描述

我是Airflow新手,为减少维护成本,尝试通过模板生成多个代码相似的DAG而非逐个创建。简单场景可正常运行,但当DAG中包含dag_run.conf、var.val.get这类Airflow模板字段时,外层Jinja会尝试渲染它们,抛出'dag_run'未定义的错误。

报错信息

Traceback (most recent call last):
  File "C:\Users\user7\Git\airflow-test\airflow_new_dag_generator.py", line 17, in <module>
    output = template.render(
  File "C:\Users\user7\AppData\Local\Programs\Python\Python39\lib\site-packages\jinja2\environment.py", line 1090, in render
    self.environment.handle_exception()
  File "C:\Users\user7\AppData\Local\Programs\Python\Python39\lib\site-packages\jinja2\environment.py", line 832, in handle_exception
    reraise(*rewrite_traceback_stack(source=source))
  File "C:\Users\user7\AppData\Local\Programs\Python\Python39\lib\site-packages\jinja2\_compat.py", line 28, in reraise
    raise value.with_traceback(tb)
  File "C:\Users\user7\Git\airflow-test\templates\airflow_new_dag_template.py", line 41, in top-level template code
    bash_command="echo {{ dag_run.conf.get('some_number')}}"
  File "C:\Users\user7\AppData\Local\Programs\Python\Python39\lib\site-packages\jinja2\environment.py", line 471, in getattr
    return getattr(obj, attribute)
jinja2.exceptions.UndefinedError: 'dag_run' is undefined

涉及代码

airflow_test_dag_template.py

from airflow import DAG
from airflow.operators.dummy import DummyOperator
from airflow.operators.bash import BashOperator
from datetime import datetime, timedelta
import os

DAG_ID: str = os.path.basename(__file__).replace(".py", "")
CITY = "{{city}}"
STATE = "{{state}}"

DEFAULT_ARGS = {
    'owner': 'airflow_test',
    'depends_on_past': False,
    'email': ['airflow@example.com'],
    'email_on_failure': True,
    'email_on_retry': False,
}

with DAG(
        dag_id=DAG_ID,
        default_args=DEFAULT_ARGS,
        dagrun_timeout=timedelta(hours=12),
        start_date=datetime(2023, 1, 1),
        catchup=False,
        schedule_interval=None,
        tags=['test']
) as dag:
    # Defining operators

    t1 = BashOperator(
        task_id="t1",
        bash_command=f"echo INFO ==> City : {CITY}, State: {STATE}"
    )

    t2 = BashOperator(
        task_id="t2",
        bash_command="echo {{ dag_run.conf.get('some_number')}}"
    )

    # Execution flow for operators
    t1 >> t2

airflow_test_dag_generator.py

from pathlib import Path
from jinja2 import Environment, FileSystemLoader

file_loader = FileSystemLoader(Path(__file__).parent)
env = Environment(loader=file_loader)
dags_folder = 'C:/Users/user7/Git/airflow-test/dags'

template = env.get_template('templates/airflow_test_dag_template.py')

city_list = ['brooklyn', 'queens']

state = 'NY'

for city in city_list:
    print(f"Generating dag for {city}...")
    file_name = f"airflow_test_dag_{city}.py"
    output = template.render(
        city=city,
        state=state
    )

    with open(dags_folder + '/' + file_name, "w") as f:
        f.write(output)
    print(f"DAG file saved under {file_name}")

解决方案

针对外层Jinja渲染Airflow模板字段的冲突问题,提供三种可行解决方法:

方法1:用Jinja raw标签保护Airflow模板内容

在模板中,将Airflow的{{ ... }}语法用{% raw %}和{% endraw %}包裹,让外层Jinja跳过这部分内容的渲染,直接保留原字符串给Airflow使用。

修改后的t2算子示例:

t2 = BashOperator(
    task_id="t2",
    bash_command="echo {% raw %}{{ dag_run.conf.get('some_number')}}{% endraw %}"
)

方法2:转义Airflow模板的花括号

把Airflow模板中的{{替换为{{ '{{' }},}}替换为{{ '}}' }},外层Jinja渲染后会生成原始的Airflow模板语法。

修改后的t2算子示例:

t2 = BashOperator(
    task_id="t2",
    bash_command="echo {{ '{{ dag_run.conf.get(''some_number'') }}' }}"
)

注意:Python字符串中的单引号需要用双单引号转义。

方法3:改用Airflow原生循环生成DAG(推荐)

不需要外部Jinja模板,直接在Python代码中通过循环生成DAG,彻底避免两层模板引擎冲突,同时符合Airflow最佳实践。

示例代码:

from airflow import DAG
from airflow.operators.bash import BashOperator
from datetime import datetime, timedelta

DEFAULT_ARGS = {
    'owner': 'airflow_test',
    'depends_on_past': False,
    'email': ['airflow@example.com'],
    'email_on_failure': True,
    'email_on_retry': False,
}

city_list = ['brooklyn', 'queens']
state = 'NY'

for city in city_list:
    dag_id = f"airflow_test_dag_{city}"
    with DAG(
        dag_id=dag_id,
        default_args=DEFAULT_ARGS,
        dagrun_timeout=timedelta(hours=12),
        start_date=datetime(2023, 1, 1),
        catchup=False,
        schedule_interval=None,
        tags=['test']
    ) as dag:
        t1 = BashOperator(
            task_id="t1",
            bash_command=f"echo INFO ==> City : {city}, State: {state}"
        )

        t2 = BashOperator(
            task_id="t2",
            bash_command="echo {{ dag_run.conf.get('some_number')}}"
        )

        t1 >> t2
    # 将生成的DAG对象存入全局变量,Airflow才能识别加载
    globals()[dag_id] = dag

内容的提问来源于stack exchange,提问作者ACL

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 18:40:56