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

如何通过Databricks记录ADF管道成功运行信息至表并动态获取字段

动态获取ADF管道元数据并记录至表中

1. 在ADF管道中传递上下文参数给Databricks Notebook

在ADF的Databricks Notebook活动里,通过参数模块直接传递管道的运行元数据,所有值都可以通过ADF系统变量动态获取:

  • 管道名称:@pipeline().PipelineName
  • 管道运行ID:@pipeline().RunId
  • 管道启动时间:@pipeline().TriggerTime

操作步骤:

  • 打开管道中的Databricks Notebook活动,切换到「参数」标签
  • 添加以下键值对:
    • 键:pipeline_name,值:@pipeline().PipelineName
    • 键:run_id,值:@pipeline().RunId
    • 键:start_time,值:@pipeline().TriggerTime

2. 在Databricks Notebook中接收参数并计算运行耗时

在Notebook里读取传入的参数,计算管道运行耗时(以秒为单位),示例Python代码:

# 定义并获取ADF传递的参数
dbutils.widgets.text("pipeline_name", "", "Pipeline Name")
dbutils.widgets.text("run_id", "", "Run ID")
dbutils.widgets.text("start_time", "", "Start Time")

pipeline_name = dbutils.widgets.get("pipeline_name")
run_id = dbutils.widgets.get("run_id")
start_time_str = dbutils.widgets.get("start_time")

# 转换时间格式并计算耗时
from datetime import datetime

start_dt = datetime.fromisoformat(start_time_str.replace("Z", "+00:00"))
end_dt = datetime.utcnow()
run_duration = (end_dt - start_dt).total_seconds()

# 运行状态固定为Success(仅管道成功时才会执行此Notebook)
run_status = "Success"

3. 将元数据写入目标表

以Delta表为例,将获取到的信息写入日志表:

# 构造数据行
log_data = [(pipeline_name, run_id, run_status, 
             start_dt.strftime("%Y-%m-%d %H:%M:%S"), 
             end_dt.strftime("%Y-%m-%d %H:%M:%S"), 
             run_duration)]

# 转为DataFrame并写入表
log_df = spark.createDataFrame(
    log_data,
    schema=["pipeline_name", "run_id", "run_status", "start_time", "end_time", "run_duration_seconds"]
)

# 追加模式写入,若表不存在则自动创建
log_df.write.mode("append").saveAsTable("your_catalog.your_schema.adf_pipeline_run_logs")

注意事项

  • 确保ADF的Databricks链接拥有目标表所在存储与元数据目录的读写权限
  • TriggerTime对应管道实际启动时间,无论手动还是触发器触发都适用
  • 以Notebook执行时的UTC时间作为结束时间,足够满足常规耗时统计需求

内容的提问来源于stack exchange,提问作者Swati B

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.20 19:09:32