【问题标题】:How to wait for data factory pipeline to complete the execution?如何等待数据工厂管道完成执行?
【发布时间】:2020-11-27 21:52:11
【问题描述】:

网络核心项目。我正在使用我的 .net 核心代码调用一些 ADF(Azure 数据工厂)管道,如下所示。

 public async Task<string> RunADFPipeline(DataFactoryManagementClient dataFactoryManagementClient, Dictionary<string,object> keyValuePairs, ADFClient aDFClient, string pieplineName)
        {
            CreateRunResponse runResponse = dataFactoryManagementClient.Pipelines.CreateRunWithHttpMessagesAsync(aDFClient.ResourceGroupName, aDFClient.DataFactoryName, pieplineName, parameters: keyValuePairs).Result.Body;
            return runResponse.RunId;
        }

此管道将运行大约五分钟,并且管道会将一些数据写入 azure sql Db。现在我的要求是从 sql db 中获取数据。我有几个问题围绕着这个问题。当我的管道完成执行时,我的代码将如何知道?我在下面尝试了一些方法。

 public async Task<object> GetPipelineInfoAsync(DataFactoryManagementClient dataFactoryManagementClient, ADFClient aDFClient, string runId)
        {
            var info = await dataFactoryManagementClient.PipelineRuns.GetAsync(aDFClient.ResourceGroupName, aDFClient.DataFactoryName, runId);
            return new
            {
                RunId = info.RunId,
                PipelineName = info.PipelineName,
                InvokedBy = info.InvokedBy.Name,
                LastUpdated = info.LastUpdated,
                RunStart = info.RunStart,
                RunEnd = info.RunEnd,
                DurationInMs = info.DurationInMs,
                Status = info.Status,
                Message = info.Message
            };
        } 

通过将第一次调用收到的 RunId 传递给上述方法,我可以获得它的状态。但我不能等到执行完成。目的是 ADF 管道会将一些数据写入 db,而我需要将这些数据发送回 UI。但我不能在当前通话中等待这个。我打算为此使用 Signal R。一旦 ADF 管道完成执行,我就可以调用 GetPipelineInfoAsync 方法,如果状态为成功,那么我可以转到 db 并获取详细信息。我面临的唯一问题是在 adf 管道完成执行之前我无法阻塞主线程。有人可以帮我解决这个问题吗?任何帮助将不胜感激。谢谢

【问题讨论】:

  • 您可以通过等待异步 API 或使用同步 API 来阻塞主线程。你不能同时拥有它——等待管道成功和数据库更新,但不要阻塞主线程。您可以在运行成功时发出刷新,但在此之前您的 UI 数据将是陈旧的。

标签: c# azure asp.net-core signalr azure-data-factory


【解决方案1】:

您可以使用管理客户端的 query by factory 扩展 API,为管道名称和状态(以及时间跨度)传递过滤器:

var filters = new PipelineRunFilterParameters
{
    Filters =
    {
        new PipelineRunQueryFilter
        {
            Operand = "PipelineName",
            OperatorProperty = "Equals",
            Values = new List<string> { "my_pipeline", "my_other_pipeline" }
        },
        new PipelineRunQueryFilter
        {
            Operand = "Status",
            OperatorProperty = "Equals",
            Values = new List<string> { "InProgress", "Queued" }
        }
    },
    LastUpdatedBefore = DateTime.Now,
    LastUpdatedAfter = (DateTime.Now.AddDays(-1))
};

while ((_adfClient.PipelineRuns.QueryByFactory(_resourceGroup, _dataFactoryName, filters).Value.Count) >= _maxConcurrentPipelines)
{
    Task.Delay(1000).Wait();
}

您还可以在创建运行和此 API 时使用客户端的响应 ID 在单个管道上等待:

while ((pipelineRun = _adfClient.PipelineRuns.Get(_resourceGroup, _dataFactoryName, response.RunId)).Status == "InProgress" || pipelineRun.Status == "Queued")
{
    Task.Delay(1000).Wait();
}

【讨论】:

    猜你喜欢
    • 2020-11-29
    • 1970-01-01
    • 2016-10-16
    • 1970-01-01
    • 2015-09-16
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多