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

如何使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.12 05:16:46