如何从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
相关产品推荐
相关产品推荐

