开启记录不兼容行后,如何用Azure Function获取ADF复制活动的skippedRowCount?
获取ADF复制活动的skippedRowCount值
当然可以获取这个值!下面是几种实用的方法,完美适配你提到的Azure Function + Azure SDK的场景:
方法一:用Azure Data Factory SDK直接查询活动运行详情
你可以通过ADF的.NET/Python SDK,主动查询指定管道运行的复制活动输出,从中提取skippedRowCount:
.NET SDK示例
- 先安装ADF NuGet包:
Install-Package Microsoft.Azure.Management.DataFactory - 编写查询代码:
using Microsoft.Azure.Management.DataFactory; using Microsoft.Azure.Management.DataFactory.Models; using Microsoft.Rest; using System.Linq; // 初始化ADF客户端 var credentials = new TokenCredentials("<你的ADF访问令牌>"); var adfClient = new DataFactoryManagementClient(credentials) { SubscriptionId = "<你的订阅ID>" }; // 查询最近1小时内的管道运行 var pipelineRuns = adfClient.PipelineRuns.ListByFactory("<资源组名>", "<ADF工厂名>", new RunFilterParameters(DateTime.UtcNow.AddHours(-1), DateTime.UtcNow)); // 定位到目标管道的运行实例 var targetRun = pipelineRuns.First(r => r.PipelineName == "<你的管道名称>"); // 获取该管道运行下的所有活动运行记录 var activityRuns = adfClient.ActivityRuns.ListByPipelineRun("<资源组名>", "<ADF工厂名>", targetRun.RunId, DateTime.UtcNow.AddHours(-1), DateTime.UtcNow); // 找到复制活动并提取跳过行数 var copyActivity = activityRuns.First(a => a.ActivityName == "<你的复制活动名称>"); if (copyActivity.Output != null && copyActivity.Output.TryGetValue("skippedRowCount", out var skippedCount)) { Console.WriteLine($"本次复制跳过的行数: {skippedCount}"); // 这里可以把值存入数据库或做其他处理 }
Python SDK示例
- 安装SDK依赖:
pip install azure-mgmt-datafactory azure-identity - 代码实现:
from azure.identity import DefaultAzureCredential from azure.mgmt.datafactory import DataFactoryManagementClient from datetime import datetime, timedelta # 初始化ADF客户端 credential = DefaultAzureCredential() adf_client = DataFactoryManagementClient(credential, "<你的订阅ID>") # 设置查询时间范围(最近1小时) start_time = datetime.utcnow() - timedelta(hours=1) end_time = datetime.utcnow() # 获取管道运行列表 pipeline_runs = adf_client.pipeline_runs.list_by_factory( "<资源组名>", "<ADF工厂名>", filter_parameters={"lastUpdatedAfter": start_time, "lastUpdatedBefore": end_time} ) # 定位目标管道运行 target_run = next(run for run in pipeline_runs if run.pipeline_name == "<你的管道名称>") # 获取活动运行记录 activity_runs = adf_client.activity_runs.list_by_pipeline_run( "<资源组名>", "<ADF工厂名>", target_run.run_id, start_time, end_time ) # 提取复制活动的跳过行数 copy_activity = next(act for act in activity_runs if act.activity_name == "<你的复制活动名称>") if "skippedRowCount" in copy_activity.output: skipped_count = copy_activity.output["skippedRowCount"] print(f"本次复制跳过的行数: {skipped_count}")
方法二:用Azure Function监听ADF运行事件实时获取
如果需要实时捕获skippedRowCount,可以通过Event Grid触发Azure Function,当ADF复制活动完成时自动提取值:
- 在ADF的事件订阅中,将管道/活动运行完成的事件发送到Event Grid主题,或者直接绑定到Azure Function。
- 编写Function处理事件(以C#为例):
using Microsoft.Azure.WebJobs; using Microsoft.Azure.WebJobs.Extensions.EventGrid; using Microsoft.Extensions.Logging; using Azure.Messaging.EventGrid; using System.Text.Json; using System.Collections.Generic; public static void Run([EventGridTrigger] EventGridEvent eventGridEvent, ILogger log) { log.LogInformation("收到ADF运行事件: {EventId}", eventGridEvent.Id); // 解析事件中的活动输出数据 if (eventGridEvent.Data != null) { var eventData = JsonSerializer.Deserialize<Dictionary<string, object>>(eventGridEvent.Data.ToString()); if (eventData.TryGetValue("output", out var outputObj)) { var output = JsonSerializer.Deserialize<Dictionary<string, object>>(outputObj.ToString()); if (output.TryGetValue("activities", out var activitiesObj)) { var activities = JsonSerializer.Deserialize<List<Dictionary<string, object>>>(activitiesObj.ToString()); var copyActivity = activities.FirstOrDefault(a => a["name"].ToString() == "<你的复制活动名称>"); if (copyActivity != null && copyActivity.TryGetValue("output", out var actOutput)) { var actOutputDict = JsonSerializer.Deserialize<Dictionary<string, object>>(actOutput.ToString()); if (actOutputDict.TryGetValue("skippedRowCount", out var skippedCount)) { log.LogInformation("复制活动跳过行数: {Count}", skippedCount); // 这里可以将值写入SQL DW或其他存储 } } } } } }
关键注意事项
- 确保你的服务主体(Function或SDK使用的身份)拥有ADF的
Reader或Data Factory Contributor权限,否则无法查询运行详情。 - 只有当复制活动启用了记录不兼容行且确实有跳过的行时,
skippedRowCount才会出现在输出中,代码里要做好空值判断。 - 实时场景优先选Event Grid触发Function,事后统计或批量查询则用SDK更灵活。
内容的提问来源于stack exchange,提问作者phương hà
相关产品推荐
相关产品推荐

