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

Azure Synapse本地触发Pipeline时参数传递为空的问题求助

问题:本地调用Azure Synapse Pipeline传递参数至Notebook为空

我在本地调用Azure Synapse中的Pipeline并传递输入参数给Notebook时遇到问题。在Synapse Studio中调试该Pipeline时参数可正常使用,但通过本地代码触发后,file_name参数的值为空。Notebook中的参数定义已多次确认无误,问题出在本地参数传递环节。

以下是我的本地代码:

import json
import time
import requests
from azure.identity import DefaultAzureCredential

# 配置参数
workspace_name = NAME
spark_pool_name = NAME
notebook_name = NAME
file_name = NAME
api_version = "2021-06-01-preview"

# Pipeline名称
pipeline_name = NAME

# 身份验证
credential = DefaultAzureCredential()
token = credential.get_token("https://dev.azuresynapse.net/.default")

# Synapse Pipeline API基础URL
base_url = f"https://{workspace_name}.dev.azuresynapse.net"
create_pipeline_url = f"{base_url}/pipelines/{pipeline_name}?api-version={api_version}"

# 请求头
headers = {
    "Authorization": f"Bearer {token.token}",
    "Content-Type": "application/json"
}

# 匹配Azure Studio中的参数结构
pipeline_payload = {
    "properties": {
        "activities": [
            {
                "name": "NotebookActivity",
                "type": "SynapseNotebook",
                "dependsOn": [],
                "policy": {
                    "timeout": "7.00:00:00",
                    "retry": 0,
                    "retryIntervalInSeconds": 30,
                    "secureOutput": False,
                    "secureInput": False
                },
                "userProperties": [],
                "typeProperties": {
                    "notebook": {
                        "referenceName": notebook_name,
                        "type": "NotebookReference"
                    },
                    "parameters": {
                        "file_name": {
                            "value": {
                                "value": "@pipeline().parameters.file_name",
                                "type": "Expression"
                            },
                            "type": "string"
                        }
                    },
                    "sparkPool": {
                        "referenceName": spark_pool_name,
                        "type": "BigDataPoolReference"
                    },
                    "conf": {
                        "spark.dynamicAllocation.enabled": None,
                        "spark.dynamicAllocation.minExecutors": None,
                        "spark.dynamicAllocation.maxExecutors": None
                    },
                    "numExecutors": None
                }
            }
        ],
        "parameters": {
            "file_name": {
                "type": "String"
            }
        }
    }
}

try:
    # 创建Pipeline
    print(f"创建Pipeline '{pipeline_name}'...")
    create_response = requests.put(create_pipeline_url, json=pipeline_payload, headers=headers)
    
    if create_response.status_code not in [200, 201, 202]:
        print(f"创建Pipeline失败: {create_response.status_code}")
        print(create_response.text)
        exit(1)
    
    print(f"Pipeline '{pipeline_name}'创建成功!")
    
    # 运行Pipeline并传递参数
    print(f"启动Pipeline '{pipeline_name}'...")
    run_pipeline_url = f"{base_url}/pipelines/{pipeline_name}/createRun?api-version={api_version}"
    
    # 参数格式
    run_payload = {
        "parameters": {
            "file_name": file_name
        }
    }
    
    run_response = requests.post(run_pipeline_url, json=run_payload, headers=headers)

    # DEBUG信息
    print("传递的参数:", run_payload)
    print("API响应:", run_response.json())

    if run_response.status_code not in [200, 201, 202]:
        print(f"运行Pipeline失败: {run_response.status_code}")
        print(run_response.text)
        exit(1)
    
    run_id = run_response.json().get("runId")
    print(f"Pipeline运行ID: {run_id}")
    
    # 监控Pipeline运行状态
    monitor_url = f"{base_url}/pipelineruns/{run_id}?api-version={api_version}"
    print(f"可通过以下URL监控状态: {monitor_url}")
    
    # 每30秒检查一次状态
    print("监控Pipeline运行状态...")
    max_checks = 20
    current_check = 0
    
    while current_check < max_checks:
        status_response = requests.get(monitor_url, headers=headers)
        
        if status_response.status_code == 200:
            status = status_response.json().get("status")
            print(f"当前状态: {status}")
            
            if status in ["Succeeded", "Failed", "Cancelled"]:
                # 若失败则输出错误信息
                if status == "Failed" and "error" in status_response.json():
                    print(f'Pipeline错误:')
                break
        else:
            print(f"获取状态失败: {status_response.status_code}")
            print(status_response.text)
        
        time.sleep(30)
        current_check += 1
    
    if current_check >= max_checks:
        print("监控超时,请手动检查Pipeline状态。")

except Exception as e:
    print(f"发生错误: {str(e)}")

问题分析与解决方案

核心问题

代码中创建Pipeline的pipeline_payload里,Notebook参数的表达式结构错误:file_name参数嵌套了两层value和type,导致Synapse无法正确解析@pipeline().parameters.file_name表达式,最终参数值为空。

修正步骤

  1. 修改Notebook参数的表达式结构
    将pipeline_payload中typeProperties.parameters部分的file_name定义修改为:

    "parameters": {
        "file_name": {
            "value": "@pipeline().parameters.file_name",
            "type": "Expression"
        }
    }
    

    移除多余的嵌套value层,直接使用表达式字符串作为value值。

  2. 可选优化:避免重复创建Pipeline
    如果Pipeline已在Synapse Studio中配置完成,无需每次运行代码都通过requests.put重新创建/覆盖Pipeline。可删除创建Pipeline的代码块,仅保留触发运行的逻辑(run_pipeline_url相关部分),减少结构错误风险。

验证

修改后重新运行代码,触发Pipeline时传递的file_name参数将正确传入Notebook,可在Synapse Studio的Pipeline运行日志中查看参数值是否正常。


内容的提问来源于stack exchange,提问作者Olimpia Dagnino

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 09:49:51