如何编写Databricks Notebook读取指定作业统计数据并写入Snowflake表
Databricks指定作业统计数据同步Snowflake实现代码
前置要求
- 拥有Databricks工作区Jobs API访问权限,以及目标作业的查看权限
- 已在Databricks Secrets中存储Databricks API Token、Snowflake连接认证信息,避免硬编码敏感信息
- Snowflake侧已提前创建目标表,建表语句参考:
CREATE OR REPLACE TABLE <你的库名>.<你的Schema名>.databricks_job_stats ( job_name STRING, start_time TIMESTAMP, end_time TIMESTAMP, status STRING );
完整Notebook代码
步骤1:导入依赖并配置参数
import requests import json from pyspark.sql.functions import col, from_unixtime # -------------------------- 需替换为你实际的配置 -------------------------- # 要提取的两个目标作业ID TARGET_JOB_IDS = ["<你的作业1ID>", "<你的作业2ID>"] # Databricks工作区配置 DATABRICKS_HOST = "https://<你的Databricks工作区域名>" DATABRICKS_TOKEN = dbutils.secrets.get(scope="<你的secrets scope名>", key="<databricks token的key名>") # Snowflake连接配置 SNOWFLAKE_OPTIONS = { "sfUrl": "<你的Snowflake账号URL>", "sfUser": dbutils.secrets.get(scope="<你的secrets scope名>", key="<snowflake用户名的key名>"), "sfPassword": dbutils.secrets.get(scope="<你的secrets scope名>", key="<snowflake密码的key名>"), "sfDatabase": "<你的Snowflake数据库名>", "sfSchema": "<你的Snowflake Schema名>", "sfWarehouse": "<你的Snowflake计算仓库名>", "sfRole": "<你的Snowflake角色名>" } TARGET_TABLE = "databricks_job_stats" # ------------------------------------------------------------------------
步骤2:调用Jobs API拉取作业运行数据
headers = { "Authorization": f"Bearer {DATABRICKS_TOKEN}", "Content-Type": "application/json" } all_runs = [] for job_id in TARGET_JOB_IDS: # 调用List Job Runs接口,可根据需要调整limit参数拉取更多历史运行记录 url = f"{DATABRICKS_HOST}/api/2.1/jobs/runs/list?job_id={job_id}&limit=100" response = requests.get(url, headers=headers) response.raise_for_status() runs = response.json().get("runs", []) all_runs.extend(runs)
步骤3:数据清洗,保留需要的字段
# 转换为DataFrame并处理字段 processed_data = [] for run in all_runs: processed_data.append({ "job_name": run.get("job_name", ""), "start_time": run.get("start_time", 0)/1000, # API返回毫秒时间戳,转秒级 "end_time": run.get("end_time", 0)/1000 if run.get("end_time") else None, "status": run.get("state", {}).get("result_state", "RUNNING") }) # 生成Spark DataFrame并转换时间格式 df = spark.createDataFrame(processed_data) \ .withColumn("start_time", from_unixtime(col("start_time")).cast("timestamp")) \ .withColumn("end_time", from_unixtime(col("end_time")).cast("timestamp")) \ .select("job_name", "start_time", "end_time", "status") # 可执行display(df)验证数据正确性 # display(df)
步骤4:写入Snowflake表
# 写入模式可按需调整:append为追加,overwrite为覆盖整张表 df.write.format("net.snowflake.spark.snowflake") \ .options(**SNOWFLAKE_OPTIONS) \ .option("dbtable", TARGET_TABLE) \ .mode("append") \ .save()
注意事项
- Databricks作业ID可在作业详情页URL中获取,格式为
/jobs/<job_id>/runs中的数字部分 - 若需要定时同步,可直接将当前Notebook配置为Databricks定时作业,按业务需求设置调度频率
- 可按需调整API的limit参数、增加时间过滤条件,仅拉取特定时间段内的作业运行记录
- 建议给Snowflake目标表增加作业运行ID作为唯一主键,搭配MERGE逻辑避免重复写入数据
内容的提问来源于stack exchange,提问作者deg
相关产品推荐
相关产品推荐

