如何通过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
相关产品推荐
相关产品推荐

