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

如何从Synapse Notebook调用带参数的Synapse Pipeline?

在Synapse Notebook中调用Synapse Pipeline并传递参数的实现方法

完全可以从Synapse Notebook中调用Synapse Pipeline执行,并且支持传递自定义参数,下面是具体的实现方案:

核心实现方式

借助Azure官方Python SDK(azure-mgmt-datafactory)结合身份认证组件(azure-identity),可直接在Notebook中发起Pipeline运行请求,同时传递预设参数。

步骤1:安装依赖包

如果Notebook环境未安装所需SDK,先执行以下命令:

!pip install azure-identity azure-mgmt-datafactory

步骤2:编写调用代码

以下是完整可运行示例,包含身份认证、参数传递和Pipeline触发逻辑:

from azure.identity import DefaultAzureCredential
from azure.mgmt.datafactory import DataFactoryManagementClient
from azure.mgmt.datafactory.models import CreateRunParameters

# 替换为你的资源信息
SUBSCRIPTION_ID = "<你的订阅ID>"
RESOURCE_GROUP_NAME = "<你的资源组名称>"
DATA_FACTORY_NAME = "<你的Synapse工作区对应Data Factory名称>"
PIPELINE_NAME = "<要调用的Pipeline名称>"

# 初始化身份认证(Synapse托管身份自动生效)
credential = DefaultAzureCredential()
adf_client = DataFactoryManagementClient(credential, SUBSCRIPTION_ID)

# 定义要传递的参数(需与Pipeline中配置的参数名完全匹配)
pipeline_parameters = {
    "input_table_name": "sales_data_2024",
    "output_container": "processed-data",
    "batch_id": 1001
}

# 触发Pipeline运行
run_response = adf_client.pipelines.create_run(
    resource_group_name=RESOURCE_GROUP_NAME,
    factory_name=DATA_FACTORY_NAME,
    pipeline_name=PIPELINE_NAME,
    parameters=CreateRunParameters(parameters=pipeline_parameters)
)

# 输出运行ID,用于后续查询状态
print(f"Pipeline运行已触发,运行ID: {run_response.run_id}")

# 可选:查询运行状态
run_details = adf_client.pipeline_runs.get(
    resource_group_name=RESOURCE_GROUP_NAME,
    factory_name=DATA_FACTORY_NAME,
    run_id=run_response.run_id
)
print(f"当前运行状态: {run_details.status}")

关键注意事项

  • 身份权限:确保Synapse工作区的托管身份拥有Data Factory Contributor或足够权限触发Pipeline运行
  • 参数匹配:传递的参数名称必须与Pipeline中定义的参数名完全一致,类型也要匹配(如字符串、数字等)
  • 运行状态查询:可通过run_id随时查询Pipeline的执行状态、日志等信息

替代方案:使用REST API

如果不想依赖SDK,也可直接调用Azure REST API触发Pipeline,示例如下:

import requests
from azure.identity import DefaultAzureCredential

# 获取访问令牌
credential = DefaultAzureCredential()
token = credential.get_token("https://management.azure.com/.default").token

# API请求URL
url = f"https://management.azure.com/subscriptions/{SUBSCRIPTION_ID}/resourceGroups/{RESOURCE_GROUP_NAME}/providers/Microsoft.DataFactory/factories/{DATA_FACTORY_NAME}/pipelines/{PIPELINE_NAME}/createRun?api-version=2018-06-01"

# 请求体包含参数
headers = {"Authorization": f"Bearer {token}", "Content-Type": "application/json"}
body = {"parameters": pipeline_parameters}

response = requests.post(url, headers=headers, json=body)
print(f"API响应状态码: {response.status_code}")
print(f"运行ID: {response.json()['runId']}")

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.25 16:27:37