如何在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()}条数据")
失败信息查看方式
- 在Synapse Studio的Notebook运行窗口右上角,点击Spark UI进入页面。
- 切换到Jobs页,找到对应主Notebook执行的Job。
- 点击Job进入Stages页,查看各Stage的Task状态;失败的Task会标记为红色。
- 点击失败的Task,即可查看包含子Notebook报错信息的详细日志。
内容的提问来源于stack exchange,提问作者SAH
相关产品推荐
相关产品推荐

