【问题标题】:How to create dependency between activities of Pipeline for Azure Data Factory in Python如何在 Python 中为 Azure 数据工厂的 Pipeline 活动创建依赖关系
【发布时间】:2019-02-02 14:20:45
【问题描述】:

在 Azure DataFactory 管道中,我尝试让两个 CopyActivities 按顺序运行,即第一个将数据从 blob 复制到 SQL 表,然后第二个将 SQL 表复制到另一个数据库。

我尝试了下面的代码,但结果管道不依赖于活动(从工作流程图和 JSON 在 Azure UI 中检查)。当我运行管道时,我收到如下错误消息: “ErrorResponseException:模板验证失败:'模板操作'我的第二个活动 nameScope' 在第 '1' 行和列 '22521' 的 'runAfter' 属性包含不存在的操作。balababla ....” em>

在 Azure UI 中手动添加依赖项后,我可以成功运行管道。

如果有人能指出示例代码 (Python/C#/Powershell) 或文档,我将不胜感激。 我的 Python 代码:

    def createDataFactoryRectStage(self,  
                               aPipelineName, aActivityStageName, aActivityAcquireName,
                               aRectFileName, aRectDSName,
                               aStageTableName, aStageDSName,
                               aAcquireTableName, aAcquireDSName):
    adf_client = self.__getAdfClient()

    ds_blob = AzureBlobDataset(linked_service_name = LinkedServiceReference(AZURE_DATAFACTORY_LS_BLOB_RECT), 
                               folder_path=PRJ_AZURE_BLOB_PATH_RECT, 
                               file_name = aRectFileName,
                               format = {"type": "TextFormat",
                                         "columnDelimiter": ",",
                                         "rowDelimiter": "",
                                         "nullValue": "\\N",
                                         "treatEmptyAsNull": "true",
                                         "firstRowAsHeader": "true",
                                         "quoteChar": "\"",})    
    adf_client.datasets.create_or_update(AZURE_RESOURCE_GROUP, AZURE_DATAFACTORY, aRectDSName, ds_blob)

    ds_stage= AzureSqlTableDataset(linked_service_name = LinkedServiceReference(AZURE_DATAFACTORY_LS_SQLDB_STAGE), 
                                   table_name='[dbo].[' + aStageTableName + ']')      
    adf_client.datasets.create_or_update(AZURE_RESOURCE_GROUP, AZURE_DATAFACTORY, aStageDSName, ds_stage)

    ca_blob_to_stage = CopyActivity(aActivityStageName, 
                                    inputs=[DatasetReference(aRectDSName)], 
                                    outputs=[DatasetReference(aStageDSName)], 
                                    source= BlobSource(), 
                                    sink= SqlSink(write_batch_size = AZURE_SQL_WRITE_BATCH_SIZE))

    ds_acquire= AzureSqlTableDataset(linked_service_name = LinkedServiceReference(AZURE_DATAFACTORY_LS_SQLDB_ACQUIRE), 
                                     table_name='[dbo].[' + aAcquireTableName + ']')      
    adf_client.datasets.create_or_update(AZURE_RESOURCE_GROUP, AZURE_DATAFACTORY, aAcquireDSName, ds_acquire)
    dep = ActivityDependency(ca_blob_to_stage, dependency_conditions =[DependencyCondition('Succeeded')])

    ca_stage_to_acquire = CopyActivity(aActivityAcquireName, 
                                       inputs=[DatasetReference(aStageDSName)], 
                                       outputs=[DatasetReference(aAcquireDSName)], 
                                       source= SqlSource(), 
                                       sink= SqlSink(write_batch_size = AZURE_SQL_WRITE_BATCH_SIZE),
                                       depends_on=[dep])

    p_obj = PipelineResource(activities=[ca_blob_to_stage, ca_stage_to_acquire], parameters={})

    return adf_client.pipelines.create_or_update(AZURE_RESOURCE_GROUP, AZURE_DATAFACTORY, aPipelineName, p_obj)

【问题讨论】:

    标签: python dependencies azure-data-factory


    【解决方案1】:

    以防万一有人像我一样遇到与这个老问题相同的问题,python 代码中有一个微妙的错误

    更改 dep 以使用活动名称,而不是对活动对象的引用使其对我有用。

    dep = ActivityDependency(aActivityStageName, dependency_conditions =[DependencyCondition('Succeeded')])
    

    【讨论】:

      【解决方案2】:

      这是C# 中的一个示例,它基本上执行Chaining activities 并在管道中按顺序链接活动。请记住,在 ADFV1 中,我们必须将一个活动的输出配置为另一个活动的输入,以链接它们并使它们相互依赖。

      管道代码 sn-p(注意 dependsOn 属性,它确保第二个活动在第一个活动完成后运行,它成功运行) -

      static PipelineResource PipelineDefinition(DataFactoryManagementClient client) {
       Console.WriteLine("Creating pipeline " + pipelineName + "...");
       PipelineResource resource = new PipelineResource {
         Activities = new List < Activity > {
          new CopyActivity {
           Name = copyFromBlobToSQLActivity,
            Inputs = new List < DatasetReference > {
             new DatasetReference {
              ReferenceName = blobSourceDatasetName
             }
            },
            Outputs = new List<DatasetReference>
            {
             new DatasetReference {
              ReferenceName = sqlDatasetName
             }
            },
            Source = new BlobSource {},
            Sink = new SqlSink {}
          },
          new CopyActivity {
           Name = copyToSQLServerActivity,
            Inputs = new List < DatasetReference > {
             new DatasetReference {
              ReferenceName = sqlDatasetName
             }
            },
            Outputs = new List<DatasetReference>
            {
             new DatasetReference {
              ReferenceName = sqlDestinationDatasetName
             }
            },
            Source = new SqlSource {},
            Sink = new SqlSink {},
            DependsOn = new List < ActivityDependency > {
             new ActivityDependency {
              Activity = copyFromBlobToSQLActivity,
               DependencyConditions = new List < String > {
                "Succeeded"
               }
             }
            }
          }
         }
       };
       Console.WriteLine(SafeJsonConvert.SerializeObject(resource, client.SerializationSettings));
       return resource;
      }
      

      请查看 ADFV2 教程 here 以获得全面的解释和更多场景。

      【讨论】:

      • 谢谢阿布舍克。这个 C# sn-p 和我上面的 Python sn-p 完全一样,除了它的内联样式。我实际上使用了 ADFV2 Python 示例代码来构建我的 sn-p。我不知道为什么我的 Python sn-p 没有设置依赖项。我实际上分别测试了这两个 CopyActivites,它们都可以工作。你检查过你的管道的 JSON 和图表吗?
      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2020-06-06
      • 1970-01-01
      • 1970-01-01
      • 2021-07-02
      • 2020-06-11
      • 1970-01-01
      相关资源
      最近更新 更多