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

Azure Databricks循环调用Scala Notebook致Spark Driver意外停止求助

问题描述

我在Azure Databricks中有一个Python笔记本,需执行137次循环迭代。每次迭代通过dbutils.notebook.run调用另一个Scala笔记本,该Scala笔记本从MongoDB查询生成DataFrame,并通过df.createOrReplaceGlobalTempView("<<view_name>>")创建全局临时视图,供Python笔记本读取并继续处理。

Python侧核心代码:

global_temporary_database = spark.conf.get("spark.sql.globalTempDatabase")
for _ in range(137):
    dbutils.notebook.run(path="<<path_to_scala_notebook>>", timeout_seconds=600, arguments=<<current_configuration>>)
    # 恢复数据并删除全局临时视图
    df = table(f"{global_temporary_database}.<<view_name>>")
    spark.catalog.dropGlobalTempView("<<view_name>>")
    # 执行行过滤、列重命名等处理

少量迭代可正常运行,但完整执行循环时出现错误:

The spark driver has stopped unexpectedly and is restarting. Your notebook will be automatically reattached

已尝试的无效方案:

  • 添加time.sleep()增加迭代间隔
  • 每次迭代后调用spark.catalog.clearCache()

集群规格:

  • Databricks Runtime版本:9.1 LTS(含Apache Spark 3.1.2、Scala 2.12)
  • 2个Worker节点:61GB内存、8核
  • 1个Driver节点:16GB内存、4核

限制条件:因需使用Scala库处理DataFrame且需跨笔记本共享数据,无法合并两个笔记本。

解决方案

1. 升级Driver节点配置

当前Driver仅16GB内存,137次迭代中,即使删除临时视图,Driver仍可能积累未释放的元数据、中间对象或Spark作业残留状态。建议:

  • 将Driver节点规格升级至32GB内存/8核(或更高),匹配Worker节点资源量级,避免成为瓶颈。
  • 集群启动时添加Spark配置:spark.driver.memoryOverhead 8192(分配8GB堆外内存,缓解堆内存压力)。

2. 替换全局临时视图为分布式存储传递数据

全局临时视图依赖Driver维护元数据,多次创建/删除易引发内存泄漏或状态混乱。改用dbutils.notebook.exit()返回临时存储路径:

  • Scala笔记本修改:处理完DataFrame后写入DBFS临时路径,返回该路径:
import java.util.UUID
val tempPath = s"dbfs:/tmp/iter_${UUID.randomUUID().toString}.parquet"
df.write.mode("overwrite").parquet(tempPath)
dbutils.notebook.exit(tempPath)
  • Python笔记本修改:接收路径读取数据,完成后删除临时文件:
for _ in range(137):
    temp_path = dbutils.notebook.run(path="<<path_to_scala_notebook>>", timeout_seconds=600, arguments=<<current_configuration>>)
    df = spark.read.parquet(temp_path)
    dbutils.fs.rm(temp_path, recurse=True)
    # 执行行过滤、列重命名等处理

该方式将数据存储转移到分布式存储,彻底规避全局临时视图的状态依赖,降低Driver压力。

3. 强化迭代后的资源清理

除clearCache()外,添加以下操作强制释放Driver资源:

import gc
# 显式触发Python垃圾回收
gc.collect()
# 重置Catalog状态(仅无并发操作时使用)
spark.sessionState.catalog.reset()
# 若DataFrame有缓存,显式取消持久化
df.unpersist()

4. 拆分迭代批次

将137次迭代拆分为多个小批次(如每20次为一批),每个批次执行完后通过作业调度重启Python笔记本会话,彻底清除Driver积累的状态。需提前规划批次间的状态保存逻辑(如将中间结果写入持久化存储)。

5. 排查Scala笔记本资源泄漏

确认Scala笔记本无未释放的资源:

  • 末尾添加spark.catalog.clearCache()
  • 显式关闭MongoDB客户端连接(若使用自定义连接池)
  • 避免保留全局变量或静态对象,确保每次执行后临时对象被回收

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 19:50:19