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

如何用C#代码检查Blob触发器触发的ADF管道运行状态

检查Blob事件触发的Azure Data Factory管道运行状态(C#实现)

1. 解决PipelineRunFilterParameters命名空间缺失问题

PipelineRunFilterParameters类属于Microsoft.Azure.Management.DataFactory.Models命名空间,确保你已安装Microsoft.Azure.Management.DataFactory NuGet包,可通过NuGet包管理器或以下命令安装:

Install-Package Microsoft.Azure.Management.DataFactory

2. 完整实现流程

步骤1:配置身份验证

使用服务主体身份验证获取DataFactoryManagementClient实例,需提前准备租户ID、客户端ID、客户端密钥、订阅ID:

using Microsoft.Azure.Management.DataFactory;
using Microsoft.Azure.Management.DataFactory.Models;
using Microsoft.Identity.Client;
using System;
using System.Linq;
using System.Threading.Tasks;

var tenantId = "你的租户ID";
var clientId = "你的客户端ID";
var clientSecret = "你的客户端密钥";
var subscriptionId = "你的订阅ID";
var resourceGroupName = "资源组名称";
var dataFactoryName = "Data Factory名称";
var pipelineName = "目标管道名称";
var blobTriggerName = "你的Blob事件触发器名称";

// 创建身份凭据
var credential = new ClientSecretCredential(tenantId, clientId, clientSecret);

// 初始化DataFactory客户端
var dataFactoryClient = new DataFactoryManagementClient(credential)
{
    SubscriptionId = subscriptionId
};

步骤2:筛选Blob事件触发的管道运行记录

通过PipelineRunFilterParameters设置筛选条件,指定时间范围、管道名称、触发类型,定位目标运行记录:

// 设置查询时间范围(按需调整,比如最近2小时)
var startTime = DateTime.UtcNow.AddHours(-2);
var endTime = DateTime.UtcNow;

// 构建筛选参数
var filterParams = new PipelineRunFilterParameters(
    startTime: startTime,
    endTime: endTime,
    pipelineName: pipelineName,
    triggerType: "BlobEventTrigger"
);

// 查询符合条件的管道运行记录
var pipelineRuns = await dataFactoryClient.PipelineRuns.QueryByFactoryAsync(
    resourceGroupName: resourceGroupName,
    factoryName: dataFactoryName,
    filterParameters: filterParams
);

步骤3:获取并检查运行状态

遍历查询结果,可匹配触发器名称或通过触发输入信息定位特定Blob触发的运行记录:

// 遍历所有符合条件的运行记录
foreach (var run in pipelineRuns.Value)
{
    if (run.TriggerName == blobTriggerName)
    {
        Console.WriteLine($"管道运行ID: {run.RunId}");
        Console.WriteLine($"运行状态: {run.Status}");
        Console.WriteLine($"开始时间: {run.StartTime}");
        Console.WriteLine($"结束时间: {run.EndTime}");
    }
}

// 获取最新一次运行记录的状态
var latestRun = pipelineRuns.Value.OrderByDescending(r => r.StartTime).FirstOrDefault();
if (latestRun != null)
{
    Console.WriteLine($"管道最新运行状态: {latestRun.Status}");
}

关键补充

  • 若需精准匹配某一特定Blob文件触发的运行,可通过run.TriggerRunId调用dataFactoryClient.TriggerRuns.GetAsync(),获取触发器运行的输入参数,从中解析Blob文件路径进行匹配。
  • 确保使用的服务主体拥有Data Factory的Data Factory Contributor或Reader权限,避免出现权限不足的错误。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 18:15:40