如何解决Airflow DAG向Databricks Spark Python任务传参失败问题?
问题:Airflow触发Databricks任务时参数传递失败,提示InputWidgetNotDefined
任务触发成功但参数未生效,报错:com.databricks.dbutils_v1.InputWidgetNotDefined: No input widget named name is defined
预期输出为Hello <Name>!,其中name为传入参数。
Databricks任务源代码
value = dbutils.widgets.get("name") print(f"Hello {value}!")
Airflow DAG源代码
from airflow import DAG from airflow.providers.databricks.operators.databricks import DatabricksSubmitRunOperator from airflow.utils.dates import days_ago import json default_args = { 'owner': 'airflow' } with DAG('databricks_dag', start_date = days_ago(2), schedule_interval = None, default_args = default_args ) as dag: opr_run_now1 = DatabricksSubmitRunOperator( task_id = 'run_now_1', json={ "existing_cluster_id": "1234-56789-10112234", "spark_python_task": { "python_file": "/Repos/username/ace-databricks-jobs/hello.py", "parameters": json.dumps({ "name": "Peter Pan", }), } }, databricks_conn_id='conn-databricks', )
问题原因
- 参数机制不匹配:
dbutils.widgets.get读取的是Databricks输入小部件参数,而Airflow的spark_python_task.parameters传递的是Python脚本的命令行位置参数,二者不属于同一参数体系,导致脚本找不到对应小部件。 - 参数格式错误:
parameters字段要求为字符串数组,代码中传入了JSON序列化后的字典,格式不符合Databricks API要求。
解决方案
方案一:修改Databricks脚本,使用命令行参数接收
将脚本改为读取命令行参数,适配Airflow传递的参数格式:
import sys # sys.argv[0]为脚本文件名,sys.argv[1]为第一个传入的参数 value = sys.argv[1] print(f"Hello {value}!")
同时修改Airflow代码中的parameters为字符串数组:
"spark_python_task": { "python_file": "/Repos/username/ace-databricks-jobs/hello.py", "parameters": ["Peter Pan"] # 直接传递字符串数组 }
方案二:改用Notebook任务并传递输入小部件参数
如果希望保留dbutils.widgets的用法,需将任务改为Notebook类型,Airflow配置对应调整:
opr_run_now1 = DatabricksSubmitRunOperator( task_id='run_now_1', json={ "existing_cluster_id": "1234-56789-10112234", "notebook_task": { "notebook_path": "/Repos/username/ace-databricks-jobs/hello", # 替换为你的Notebook路径 "widgets": { "name": "Peter Pan" } } }, databricks_conn_id='conn-databricks', )
Notebook代码保持原有的dbutils.widgets.get逻辑即可。
内容的提问来源于stack exchange,提问作者Igor L.
相关产品推荐
相关产品推荐

