Airflow Web界面通过Trigger DAG w/ conf传递数组参数至DAG时访问失败的问题排查
问题分析与解决方案
首先咱们拆解你遇到的两个核心问题:JSON参数格式错误和DAG解析阶段访问dag_run的逻辑错误,一步步来解决:
1. 先修正触发时的JSON参数格式
你计划传递的JSON存在语法错误:数组[]里不能直接写键值对"K1" : "V1",不符合JSON规范。根据你的需求,两种正确格式如下:
- 如果要传递一组键值对,改为对象数组:
{"list_key": [{"K1": "V1"}, {"K2": "V2"}]} - 如果是普通字符串列表,直接写字符串数组:
{"list_key": ["V1", "V2"]}
2. 为什么会报NameError: name 'dag_run' is not defined?
Airflow加载DAG分为两个关键阶段:
- 解析阶段:Airflow定期扫描DAG文件,生成DAG对象,这个阶段没有运行中的DAG实例,
dag_run对象还不存在。 - 运行阶段:DAG被触发后,才会创建对应的
dag_run实例,此时才能访问dag_run.conf。
你在DAG文件的顶层代码(解析阶段执行)直接写args = list({{ dag_run.conf['list_key'] }}),不仅dag_run还未生成,而且Jinja模板语法{{ ... }}不能直接嵌入Python代码——它只能在Airflow支持模板化的Operator参数中使用。
3. 修正后的DAG代码
DataprocSubmitJobOperator的job参数默认支持模板化,我们可以直接在pyspark_job.args里用Jinja模板动态引用dag_run.conf,不需要在顶层代码提前赋值。
修正后的完整代码:
from airflow import DAG from airflow.providers.google.cloud.operators.dataproc import DataprocSubmitJobOperator from datetime import datetime import pendulum from helper import help config = help.loadJSON("config/loc") common_task_args = { 'owner': 'me', 'depends_on_past': False # 更多配置项 } dag = DAG('my-dag', default_args=common_task_args, is_paused_upon_creation=True, catchup=False, schedule_interval=None) # 直接在job定义中使用Jinja模板,动态获取dag_run.conf参数 MY_TASK = { "reference": {"project_id": "project-id"}, "placement": {"cluster_name": "cluster-name"}, "pyspark_job": { "main_python_file_uri": config["pyspark_uri"], "properties": config["spark_properties"], # 使用Jinja模板获取conf,同时设置默认值避免无参数时出错 "args": {{ dag_run.conf.get('list_key', []) }} } } my_task = DataprocSubmitJobOperator( task_id="my_task", job=MY_TASK, dag=dag )
额外注意事项
- 如果你的
list_key是对象数组,Spark的args会把每个对象转成字符串传递,需要在PySpark代码里把字符串转回JSON对象才能正常使用。 - 触发DAG时务必保证输入的JSON格式正确,否则Airflow无法解析
dag_run.conf,会导致任务失败。
内容的提问来源于stack exchange,提问作者mang4521
相关产品推荐
相关产品推荐

