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

开启记录不兼容行后,如何用Azure Function获取ADF复制活动的skippedRowCount?

获取ADF复制活动的skippedRowCount值

当然可以获取这个值!下面是几种实用的方法,完美适配你提到的Azure Function + Azure SDK的场景:

方法一:用Azure Data Factory SDK直接查询活动运行详情

你可以通过ADF的.NET/Python SDK,主动查询指定管道运行的复制活动输出,从中提取skippedRowCount:

.NET SDK示例

  1. 先安装ADF NuGet包:Install-Package Microsoft.Azure.Management.DataFactory
  2. 编写查询代码:
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示例

  1. 安装SDK依赖:pip install azure-mgmt-datafactory azure-identity
  2. 代码实现:
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复制活动完成时自动提取值:

  1. 在ADF的事件订阅中,将管道/活动运行完成的事件发送到Event Grid主题,或者直接绑定到Azure Function。
  2. 编写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à

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 06:56:55