如何让Databricks Python脚本同时获取Airflow参数与任务级参数
解决方案
方法1:通过环境变量获取Databricks任务元数据(推荐)
Databricks运行Python任务时,会自动将任务相关元数据注入到环境变量中,完全无需依赖dbutils。你可以用os.environ读取这些任务级参数,同时用sys.argv获取Airflow传递的自定义参数,两者互不干扰。
示例脚本代码:
import sys import os import json # 读取Airflow传递的JSON参数 dag_params = json.loads(sys.argv[1]) if len(sys.argv) > 1 else {} # 读取Databricks自动注入的任务元数据 job_run_id = os.environ.get("DATABRICKS_RUN_ID") job_id = os.environ.get("DATABRICKS_JOB_ID") # 合并并打印所有参数 all_params = {**dag_params, "job_run_id": job_run_id, "job_id": job_id} print("所有参数:", all_params)
方法2:Airflow中合并参数后传递(不推荐)
如果必须通过sys.argv一次性获取所有参数,需要在Airflow提交任务时,手动将自定义参数与任务元数据合并后传递。但注意:job.run_id是Databricks任务运行时生成的,Airflow提交前无法直接获取,需要额外调用Databricks API获取运行后的ID再补充参数,步骤繁琐:
- Airflow中提交任务并获取run_id:
from airflow.providers.databricks.operators.databricks import DatabricksSubmitRunOperator import json # 先提交基础任务 submit_task = DatabricksSubmitRunOperator( task_id="submit_databricks_job", databricks_conn_id="databricks_default", run_name="Airflow-Databricks-Job", spark_python_task={ "python_file": "/path/to/your/script.py", "parameters": [json.dumps({"custom_param": "test_value"})] } ) # 执行任务后获取run_id run_id = submit_task.execute(context=None)
- 调用Databricks API更新任务参数,将run_id加入后重新传递(需额外配置API权限,流程复杂)。
关键说明
- 当你在Airflow的
spark_python_task中指定parameters时,Databricks会用该参数覆盖默认的任务级参数命令行传递行为,因此默认的job.run_id不会出现在sys.argv中。 - 环境变量方式是Databricks官方支持的标准行为,脚本可独立运行(本地测试时手动设置对应环境变量即可模拟),完全符合你的业务要求。
内容的提问来源于stack exchange,提问作者Diksha Bisht
相关产品推荐
相关产品推荐

