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

如何编写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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.30 00:54:03