如何使用Python连接多个Azure Data Factory活动及配置输入输出?
用Python SDK实现Azure Data Factory多活动连接与输入输出传递
你猜的完全没错——Azure Data Factory的门户UI本质上就是个可视化的JSON编辑器,帮你生成包含活动依赖、输入输出映射的管道定义。用Python SDK实现多活动连接其实就是手动构建这些逻辑,分两部分就能搞定:
一、先搞定活动间的控制流依赖(执行顺序)
要让一个活动在另一个活动完成后执行,核心是给目标活动添加depends_on属性,指定依赖的活动和触发条件(比如成功完成、失败、无论结果)。
举个实际例子:先执行一个Copy活动,再执行Lookup活动,要求Copy成功后才启动Lookup:
from azure.mgmt.datafactory.models import ( PipelineResource, CopyActivity, LookupActivity, ActivityDependency, SuccessCondition ) # 1. 定义第一个活动:从Blob复制数据到SQL copy_activity = CopyActivity( name="CopyFromBlobToSQL", source=your_blob_source_def, # 替换成你的Blob数据源定义 sink=your_sql_sink_def # 替换成你的SQL目标定义 ) # 2. 定义第二个活动:查询SQL表数据 lookup_activity = LookupActivity( name="LookupSQLTable", source=your_sql_lookup_source_def # 替换成你的Lookup数据源定义 ) # 3. 给Lookup添加依赖:必须等Copy活动成功完成 lookup_activity.depends_on = [ ActivityDependency( activity="CopyFromBlobToSQL", dependency_conditions=[SuccessCondition.SUCCESS] ) ] # 4. 把两个活动组装成管道并创建 pipeline = PipelineResource( activities=[copy_activity, lookup_activity] ) # 省略ADF客户端初始化逻辑,直接创建管道 adf_client.pipelines.create_or_update( resource_group_name="你的资源组名", factory_name="你的ADF工厂名", pipeline_name="你的管道名", pipeline=pipeline )
如果需要依赖多个活动,或者依赖条件是“完成即可(不管成功失败)”,只需要调整dependency_conditions的值,比如用CompletionCondition.COMPLETED。
二、再处理活动间的输入输出映射(数据流传递)
如果需要把前一个活动的输出结果作为后一个活动的输入,就要用到ADF的动态内容表达式——在Python里就是把表达式字符串直接写入目标活动的参数或配置中。
比如,把Copy活动的写入行数传给后续的存储过程活动作为参数:
from azure.mgmt.datafactory.models import StoredProcedureActivity, StoredProcedureParameter # 1. 定义Copy活动(确保它能返回输出结果) copy_activity = CopyActivity( name="CopyData", source=your_blob_source_def, sink=your_sql_sink_def ) # 2. 定义存储过程活动,引用Copy活动的输出 sp_activity = StoredProcedureActivity( name="LogCopyResult", stored_procedure_name="dbo.LogCopyExecution", # 你的存储过程名 parameters={ "RowsCopied": StoredProcedureParameter( # 用ADF动态表达式引用Copy活动的输出属性 value="@activity('CopyData').output.rowsCopied", type="Int" ) }, linked_service_name=your_sql_linked_service_ref # 你的SQL链接服务引用 ) # 3. 添加依赖:等Copy完成再执行存储过程 sp_activity.depends_on = [ ActivityDependency(activity="CopyData", dependency_conditions=[SuccessCondition.SUCCESS]) ] # 4. 组装并创建管道 pipeline = PipelineResource(activities=[copy_activity, sp_activity]) adf_client.pipelines.create_or_update(...)
不同类型的活动输出结构不一样:比如Lookup活动的输出是@activity('LookupName').output.firstRow(单行)或@activity('LookupName').output.value(多行数组),你需要根据活动类型调整表达式。
内容的提问来源于stack exchange,提问作者pelos
相关产品推荐
相关产品推荐

