如何通过DataBricks API编写新DAG及克隆模板DAG并替换字段
使用DataBricks API操作DAG的解决方案
1. 如何使用DataBricks API编写新的DAG?
DataBricks中的DAG通过Jobs API创建,核心步骤如下:
- 生成访问凭证:在DataBricks工作区的「用户设置」中创建个人访问令牌(PAT),需授予
jobs.create权限。 - 构造任务配置Payload:Payload需包含DAG名称、任务定义、依赖关系、调度规则等信息,示例结构如下:
{ "name": "daily_etl_dag", "tasks": [ { "task_key": "extract_data", "spark_python_task": { "python_file_path": "/Workspace/etl/scripts/extract.py" }, "depends_on": [] }, { "task_key": "transform_data", "spark_python_task": { "python_file_path": "/Workspace/etl/scripts/transform.py" }, "depends_on": [{"task_key": "extract_data"}] }, { "task_key": "load_data", "spark_sql_task": { "query": "INSERT INTO prod_db.fact_sales SELECT * FROM staging.fact_sales" }, "depends_on": [{"task_key": "transform_data"}] } ], "schedule": { "quartz_cron_expression": "0 0 2 * * ?", "timezone_id": "UTC" }, "max_concurrent_runs": 1 }
- 调用创建API:发送POST请求到
https://<你的Databricks实例域名>/api/2.1/jobs/create,请求头携带Authorization: Bearer <你的PAT令牌>和Content-Type: application/json,将上述Payload作为请求体发送。
2. 如何通过Python调用DataBricks API克隆模板DAG并替换字段?
以下是完整的实现流程,确保保留原DAG的所有配置(依赖、调度、资源等),仅替换指定字段:
步骤说明
- 拉取模板DAG的完整配置:调用Jobs API的
get接口获取模板任务的所有参数。 - 批量替换目标字段:遍历配置中的所有层级,将包含
template的标识替换为stage。 - 提交新DAG配置:将修改后的配置通过
create接口生成新任务。
Python代码示例
import requests # 基础配置 DATABRICKS_WORKSPACE = "your-workspace-url" PAT_TOKEN = "your-personal-access-token" TEMPLATE_JOB_ID = 456 # 模板DAG的Job ID REPLACE_MAP = {"template": "stage"} def fetch_template_job(job_id): """拉取模板DAG的配置""" url = f"https://{DATABRICKS_WORKSPACE}/api/2.1/jobs/get" headers = {"Authorization": f"Bearer {PAT_TOKEN}"} params = {"job_id": job_id} response = requests.get(url, headers=headers, params=params) response.raise_for_status() return response.json() def replace_config_values(config, replace_map): """递归替换配置中的指定字段""" if isinstance(config, dict): for key, value in config.items(): if isinstance(value, (dict, list)): config[key] = replace_config_values(value, replace_map) elif isinstance(value, str): for old_str, new_str in replace_map.items(): config[key] = value.replace(old_str, new_str) # 移除原Job ID,避免创建冲突 config.pop("job_id", None) return config elif isinstance(config, list): return [replace_config_values(item, replace_map) for item in config] else: return config def create_cloned_job(config): """创建新的克隆DAG""" url = f"https://{DATABRICKS_WORKSPACE}/api/2.1/jobs/create" headers = {"Authorization": f"Bearer {PAT_TOKEN}", "Content-Type": "application/json"} response = requests.post(url, headers=headers, json=config) response.raise_for_status() return response.json() if __name__ == "__main__": # 执行克隆流程 template_config = fetch_template_job(TEMPLATE_JOB_ID) cloned_config = replace_config_values(template_config, REPLACE_MAP) result = create_cloned_job(cloned_config) print(f"克隆DAG成功,新Job ID: {result['job_id']}")
注意事项
- 确保PAT令牌拥有
jobs.get和jobs.create权限,否则会返回权限错误。 - 递归替换逻辑会覆盖所有字符串类型的字段,若模板中
template出现在其他需要保留的位置,需调整替换规则。 - 克隆后的DAG会继承原模板的所有配置,包括调度规则、资源规格、重试策略等,无需额外配置。
内容的提问来源于stack exchange,提问作者Mark McGown
相关产品推荐
相关产品推荐

