Azure Synapse中工具笔记本如何获取参数选择执行的函数?
在Azure Synapse的Functions_Utils笔记本中处理传入参数并执行对应函数
1. 获取传入的参数
在Functions_Utils笔记本里,使用mssparkutils.widgets.get()方法读取主笔记本传入的参数,核心是先拿到function_name来确定要执行的函数,其他参数按需获取。
示例代码:
# 获取指定要执行的函数名 function_name = mssparkutils.widgets.get("function_name") # 先获取通用参数,其他专属参数可在对应函数分支里读取 kv_name = mssparkutils.widgets.get("kv_name")
2. 定义对应的功能函数
提前写好各个业务函数,每个函数接收所需参数并返回对应结果。
示例代码:
def get_blob_service_client(kv_name, account_name, account_key): # 替换为实际获取Blob服务客户端的逻辑 from azure.storage.blob import BlobServiceClient blob_service_client = BlobServiceClient( account_url=f"https://{account_name}.blob.core.windows.net", credential=account_key ) return blob_service_client, account_name def get_bd_connection(kv_name, server_name, user_name, password_name, database_name): # 替换为实际获取数据库连接的逻辑 from sqlalchemy import create_engine from pyspark.sql import SparkSession # 从密钥库获取密码 password = mssparkutils.credentials.getSecret(kv_name, password_name) # 创建Spark会话和SQLAlchemy引擎 spark_session = SparkSession.builder.appName("DB_Connection").getOrCreate() sql_engine = create_engine( f"mssql+pyodbc://{user_name}:{password}@{server_name}.database.windows.net:1433/{database_name}?driver=ODBC+Driver+17+for+SQL+Server" ) return spark_session, sql_engine
3. 根据函数名分支执行并返回结果
由于mssparkutils.notebook.run()仅支持字符串类型的返回值,需要将函数返回的复杂对象序列化为JSON格式(若对象无法直接序列化,可返回其关键标识或连接信息),再通过mssparkutils.notebook.exit()传递回主笔记本。
示例代码:
import json if function_name == "get_blob_service_client": # 获取当前函数专属参数 account_name = mssparkutils.widgets.get("account_name") account_key = mssparkutils.widgets.get("account_key") # 执行函数 client, acc_name = get_blob_service_client(kv_name, account_name, account_key) # 序列化结果(若client无法直接序列化,可返回必要字符串信息) result = json.dumps({ "blob_service_client_info": str(client.account_url), "account_name": acc_name }) mssparkutils.notebook.exit(result) elif function_name == "get_bd_connection": # 获取当前函数专属参数 server_name = mssparkutils.widgets.get("server_name") user_name = mssparkutils.widgets.get("user_name") password_name = mssparkutils.widgets.get("password_name") database_name = mssparkutils.widgets.get("database_name") # 执行函数 session, engine = get_bd_connection(kv_name, server_name, user_name, password_name, database_name) # 序列化结果 result = json.dumps({ "spark_session_id": session.sparkContext.applicationId, "sql_engine_url": str(engine.url) }) mssparkutils.notebook.exit(result) else: # 处理未知函数名的情况 mssparkutils.notebook.exit(json.dumps({"error": f"未识别的函数名: {function_name}"}))
注意事项
- 参数可以按需获取:非通用参数不需要提前读取,在对应函数分支内获取即可,避免无意义的参数报错。
- 序列化限制:如果返回的对象无法序列化为JSON(如SparkSession实例),建议返回该对象的关键属性(如会话ID)或连接字符串,而非对象本身。
- 异常处理:可以在函数内添加
try-except逻辑,捕获执行错误并返回明确的错误信息。
内容的提问来源于stack exchange,提问作者David Molina
相关产品推荐
相关产品推荐

