Airflow触发Databricks依赖任务:运行时参数传递失败求助
解决方案
方法1:通过Databricks全局参数引用实现参数共享
修改Databricks工作流的任务定义,将两个任务的参数都改为引用全局作业参数,这样Airflow传入的参数会自动同步到所有引用的任务中:
"tasks": [ { "task_key": "Job1", "run_if": "ALL_SUCCESS", "python_wheel_task": { "package_name": "test", "entry_point": "test", "named_parameters": { "parm1": "{{job_parameters.parm1}}", "parm2": "{{job_parameters.parm2}}", "parm3": "{{job_parameters.parm3}}", "parm4": "{{job_parameters.parm4}}", "parm5": "{{job_parameters.parm5}}" } }, "job_cluster_key": "Job_cluster", "libraries": [ { "whl": "dbfs path" } ] }, { "task_key": "Job2", "depends_on": [{"task_key": "Job1"}], "run_if": "ALL_SUCCESS", "notebook_task": { "notebook_path": "notebook path", "base_parameters": { "parm4": "{{job_parameters.parm4}}", "parm5": "{{job_parameters.parm5}}" }, "source": "WORKSPACE" } } ]
同时调整Airflow的参数传递配置,改用全局parameters字段(而非仅针对Python Wheel任务的python_named_params):
"job_id": job_id, "parameters": { "parm1": 1, "parm2": 2, "parm3": 3, "parm4": "devl", "parm5": "yyyymm" }
方法2:在Airflow中显式覆盖Job2的参数
如果不想修改Databricks的原有作业定义,可以在Airflow触发时,通过tasks字段直接覆盖Job2的参数配置:
# 以DatabricksSubmitRunOperator为例 from airflow.providers.databricks.operators.databricks import DatabricksSubmitRunOperator submit_databricks_job = DatabricksSubmitRunOperator( task_id="trigger_databricks_workflow", databricks_conn_id="databricks_default", job_id=job_id, python_named_params={ "parm1": 1, "parm2": 2, "parm3": 3, "parm4": "devl", "parm5": "yyyymm" }, tasks=[ { "task_key": "Job2", "notebook_task": { "base_parameters": { "parm4": "devl", "parm5": "yyyymm" } } } ] )
注:tasks数组只需包含需要覆盖的任务配置,Databricks会自动合并原有作业定义与传入的配置。
方法3:通过Job1的任务输出传递参数(适合动态参数场景)
如果参数需要经过Job1处理后再传递给Job2,可以在Job1的Python Wheel代码中写入任务输出:
# Job1的Python代码 from databricks.sdk import WorkspaceClient import os w = WorkspaceClient() task_run_id = os.environ.get("DATABRICKS_TASK_RUN_ID") # 将参数写入任务输出 w.task_runs.set_output( task_run_id=task_run_id, output={"parm4": "devl", "parm5": "yyyymm"} )
然后在Job2的Notebook中读取该输出:
# Job2的Notebook代码 from databricks.sdk import WorkspaceClient w = WorkspaceClient() parent_run_id = dbutils.widgets.get("parent_run_id") # 获取Job1的任务输出 job1_run_details = w.task_runs.get(task_run_id=parent_run_id, task_key="Job1") parm4 = job1_run_details.output.get("parm4") parm5 = job1_run_details.output.get("parm5")
内容的提问来源于stack exchange,提问作者Parijat Chakraborty
相关产品推荐
相关产品推荐

