Airflow使用Jinja模板读取JSON变量遇异常及任务循环问题求助
报错原因
DAG的schedule_interval参数是在DAG解析阶段就需要确定的配置,不支持Jinja模板渲染。你直接把Jinja模板字符串'{{var.json.snflk_json.schedule_interval_snwflke_acct}}'赋值给它,Airflow会把这个字符串本身当成cron表达式去解析,自然会抛出格式错误。
解决方法
虽然你要求避免使用Variable.get(),但针对DAG级别的静态配置(比如schedule_interval),必须在解析阶段直接读取变量值。可以用deserialize_json=True直接解析JSON变量:
from airflow.models import Variable # 读取并解析JSON变量 snflk_config = Variable.get("snflk_json", deserialize_json=True) with DAG( dag_id=dag_id, default_args=default_args, # 直接用解析后的变量值 schedule_interval=snflk_config["schedule_interval_snwflke_acct"], dagrun_timeout=timedelta(hours=3), max_active_runs=1, catchup=False, params={}, tags=tags ) as dag: # 后续代码...
报错原因
你把shares赋值为Jinja模板字符串'{{var.json.snflk_json.LIST}}',这本质是一个普通字符串,不是列表。遍历这个字符串时会逐个字符处理(比如[、"、{这些),导致生成的task_id包含非法字符,触发Airflow的任务ID格式校验错误。
解决方法
同样,DAG中的任务生成逻辑是在解析阶段执行的,必须直接读取变量的列表值,不能依赖Jinja渲染:
from airflow.models import Variable # 读取并解析JSON变量 snflk_config = Variable.get("snflk_json", deserialize_json=True) with DAG( dag_id=dag_id, default_args=default_args, schedule_interval=snflk_config["schedule_interval_snwflke_acct"], dagrun_timeout=timedelta(hours=3), max_active_runs=1, catchup=False, params={}, tags=tags ) as dag: # 直接用解析后的列表 shares = snflk_config["LIST"] for s in shares: sf_tasks = SnowflakeOperator( task_id=f"{s}", snowflake_conn_id=snowflake_conn_id, sql=sqls, params={"sf_env": s}, )
补充说明
Jinja模板语法{{var.json.xxx}}主要用于任务运行阶段的动态值(比如Operator的sql参数、params参数等),不能用于DAG解析阶段的静态配置(比如schedule_interval、任务生成循环)。如果一定要完全避免Variable.get(),可以考虑使用Airflow的环境变量注入或者自定义Timetable(针对调度间隔),但实现复杂度更高,不如直接用Variable.get()高效。
内容的提问来源于stack exchange,提问作者Karthik

