You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何让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再补充参数,步骤繁琐:

  1. 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)
  1. 调用Databricks API更新任务参数,将run_id加入后重新传递(需额外配置API权限,流程复杂)。

关键说明

  • 当你在Airflow的spark_python_task中指定parameters时,Databricks会用该参数覆盖默认的任务级参数命令行传递行为,因此默认的job.run_id不会出现在sys.argv中。
  • 环境变量方式是Databricks官方支持的标准行为,脚本可独立运行(本地测试时手动设置对应环境变量即可模拟),完全符合你的业务要求。

内容的提问来源于stack exchange,提问作者Diksha Bisht

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.12 11:59:57