使用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
相关产品推荐
相关产品推荐

