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

Azure Data Factory管道审计咨询:无需修改管道如何实现表更新统计?

追踪Azure Data Factory管道运行后数据库表更新的方案

不修改现有管道的实现方案

1. 基于管道运行时间戳的表行数比对

这是最直接的无侵入方案,核心逻辑是通过管道运行时间范围,比对表在运行前后的行数差异:

  • 步骤:
    • 调用Pipeline Runs API获取目标管道的运行记录,提取精确的startTime和endTime
    • 针对需要监控的数据库表,分别获取管道启动前和管道结束后的行数快照
    • 计算两次快照的行数差值,差值不为0的表即为被更新的对象
  • 注意事项:
    • 如果表有last_modified之类的时间字段,可以只筛选该时间在管道运行范围内的行,缩小比对范围
    • 若没有时间字段,直接统计全表行数即可,适合中小规模表
  • 示例SQL(以SQL Server为例):
    -- 管道运行前行数快照(提前5分钟取数避免时间偏差)
    SELECT 
      OBJECT_NAME(object_id) AS table_name,
      SUM(row_count) AS pre_run_rows
    FROM sys.dm_db_partition_stats
    WHERE index_id IN (0,1) -- 仅统计堆表或聚集索引的行数
      AND OBJECT_NAME(object_id) IN ('目标表1','目标表2')
    GROUP BY object_id;
    
    -- 管道运行后行数快照
    SELECT 
      OBJECT_NAME(object_id) AS table_name,
      SUM(row_count) AS post_run_rows
    FROM sys.dm_db_partition_stats
    WHERE index_id IN (0,1)
      AND OBJECT_NAME(object_id) IN ('目标表1','目标表2')
    GROUP BY object_id;
    

2. 利用Azure Monitor日志关联分析

如果已经开启ADF和数据库的日志收集,可以通过日志关联直接定位被修改的表:

  • 步骤:
    • 在ADF中开启诊断设置,将管道运行日志发送到Log Analytics工作区
    • 为目标数据库开启日志收集(比如SQL Server开启Query Store,或直接将数据库诊断日志发送到Log Analytics)
    • 在Log Analytics中编写Kusto查询,关联管道运行时间范围和数据库的写入操作日志(INSERT/UPDATE/DELETE/MERGE),提取涉及的表名
  • 示例Kusto查询:
    // 获取目标管道的成功运行时间范围
    let targetPipelineRuns = AzureDiagnostics
    | where ResourceType == "DATAFACTORIES" 
      and OperationName == "PipelineRunCompleted"
      and PipelineName == "你的管道名称"
      and Status == "Succeeded"
    | project PipelineRunId, 
              StartTime=parse_json(properties).startTime, 
              EndTime=parse_json(properties).endTime;
    
    // 筛选该时间段内有写入操作的表
    AzureDiagnostics
    | where ResourceType == "SQLSERVERDATABASES" 
      and OperationName == "QueryStoreRuntimeStatistics"
      and query_text has_any ("INSERT", "UPDATE", "DELETE", "MERGE")
      and TimeGenerated between (targetPipelineRuns.StartTime .. targetPipelineRuns.EndTime)
    | parse query_text with * "INTO [" SchemaName "." TableName "]" *
    | parse query_text with * "UPDATE [" SchemaName "." TableName "]" *
    | parse query_text with * "DELETE FROM [" SchemaName "." TableName "]" *
    | distinct TableName
    

可选:修改现有管道的精准方案

如果允许对现有管道做微小调整,可以实现更精准的审计:

  • 在每个数据写入活动(Copy Data、Stored Procedure等)后,添加自定义日志活动,将目标表名、行数变化直接写入审计表或Log Analytics
  • 调用Activity Runs API,直接获取每个写入活动的输出参数(比如Copy Data活动的rowsCopied字段),结合活动配置中的目标表信息,自动记录行数变化

核心结论

完全可以在不修改现有管道的前提下完成需求,优先推荐:

  • 若数据库无完善日志体系,使用「时间戳行数比对」方案
  • 若已部署Azure Monitor日志收集,使用「日志关联分析」方案

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 20:36:08