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

如何在Synapse中并行调用带参Notebook并在Spark UI显示

实现方案

核心思路

放弃Driver端的ThreadPool,改用Spark分布式计算框架(RDD/Spark DataFrame)启动并行任务。每个Executor节点上的任务独立调用子Notebook,任务的执行日志、失败信息会被Spark捕获并展示在Spark UI的Jobs -> Stages -> Task详情中。同时通过mssparkutils.notebook.run()传递参数,子Notebook的数据库连接、DataFrame操作等逻辑可直接复用。

关键注意事项

  • 子Notebook中的资源(如数据库连接)需在每个任务中独立初始化,避免跨任务资源竞争。
  • Synapse默认支持mssparkutils在Executor节点使用,无需额外配置。
  • 传入参数需可序列化,优先使用基本数据类型组成的字典;若需传递复杂对象,可先序列化为JSON字符串,子Notebook再解析还原。
代码示例

主Notebook代码

# 1. 定义子Notebook的参数列表,每个字典对应一次调用的参数
param_list = [
    {"db_name": "sales_db", "table_name": "orders", "output_path": "abfss://container@account.dfs.core.windows.net/sales/orders"},
    {"db_name": "sales_db", "table_name": "customers", "output_path": "abfss://container@account.dfs.core.windows.net/sales/customers"},
    {"db_name": "inventory_db", "table_name": "products", "output_path": "abfss://container@account.dfs.core.windows.net/inventory/products"}
]

# 2. 将参数列表转为Spark RDD,实现分布式并行执行
rdd = spark.sparkContext.parallelize(param_list)

# 3. 定义单参数对应的执行逻辑:调用子Notebook并处理结果/异常
def run_notebook(params):
    try:
        # 调用子Notebook,设置超时时间(3600秒)并传入参数
        result = mssparkutils.notebook.run("子Notebook名称", 3600, params)
        return f"成功: {params['table_name']}, 处理结果: {result}"
    except Exception as e:
        # 捕获异常并返回失败详情,该信息会被Spark UI记录
        return f"失败: {params['table_name']}, 错误信息: {str(e)}"

# 4. 执行并行任务并收集所有结果
execution_results = rdd.map(run_notebook).collect()

# 5. 打印所有执行结果
for res in execution_results:
    print(res)

子Notebook代码

# 1. 获取主Notebook传入的参数
params = mssparkutils.notebook.getParams()
db_name = params.get("db_name")
table_name = params.get("table_name")
output_path = params.get("output_path")

# 2. 读取数据库表(示例用Synapse Spark内置表连接,也可替换为JDBC连接)
df = spark.read.table(f"{db_name}.{table_name}")

# 3. 示例DataFrame转换操作
from pyspark.sql.functions import current_date
df_transformed = df.withColumn("load_date", current_date())

# 4. 示例:写入多类型文件(Parquet + CSV)
df_transformed.write.mode("overwrite").parquet(f"{output_path}/parquet")
df_transformed.write.mode("overwrite").option("header", "true").csv(f"{output_path}/csv")

# 5. 返回执行结果(可选,主Notebook可接收该返回值)
mssparkutils.notebook.exit(f"完成处理,共{df_transformed.count()}条数据")
失败信息查看方式
  1. 在Synapse Studio的Notebook运行窗口右上角,点击Spark UI进入页面。
  2. 切换到Jobs页,找到对应主Notebook执行的Job。
  3. 点击Job进入Stages页,查看各Stage的Task状态;失败的Task会标记为红色。
  4. 点击失败的Task,即可查看包含子Notebook报错信息的详细日志。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 03:07:15