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

使用Azure Databricks Autoloader BinaryFile与foreach()遇Java堆内存溢出

问题:Azure Databricks Autoloader处理大文件时OOM异常

现象

  • 使用Azure Databricks Autoloader的BinaryFile格式结合foreach()方法实现跨位置文件复制,150MB以内的小文件运行正常,处理大文件时抛出Java堆内存不足异常
  • 单独测试shutil.copy复制20GB大文件可正常运行

错误日志

22/09/07 10:25:51 INFO FileScanRDD: Reading File path: dbfs:/mnt/somefile.csv, range: 0-1652464461, partition values: [empty row], modificationTime: 1662542176000.
22/09/07 10:25:52 ERROR Utils: Uncaught exception in thread stdout writer for /databricks/python/bin/python
java.lang.OutOfMemoryError: Java heap space
    at org.apache.spark.sql.catalyst.expressions.UnsafeRow.getBinary(UnsafeRow.java:416)
    at org.apache.spark.sql.catalyst.expressions.SpecializedGettersReader.read(SpecializedGettersReader.java:75)
    at org.apache.spark.sql.catalyst.expressions.UnsafeRow.get(UnsafeRow.java:333)
    at org.apache.spark.sql.execution.python.EvaluatePython$.toJava(EvaluatePython.scala:58)
    at org.apache.spark.sql.execution.python.PythonForeachWriter.$anonfun$inputByteIterator$1(PythonForeachWriter.scala:43)
    at org.apache.spark.sql.execution.python.PythonForeachWriter$$Lambda$1830/1643360976.apply(Unknown Source)
    at scala.collection.Iterator$$anon$10.next(Iterator.scala:461)
    at org.apache.spark.api.python.SerDeUtil$AutoBatchedPickler.next(SerDeUtil.scala:92)
    at org.apache.spark.api.python.SerDeUtil$AutoBatchedPickler.next(SerDeUtil.scala:82)
    at scala.collection.Iterator.foreach(Iterator.scala:943)
    at scala.collection.Iterator.foreach$(Iterator.scala:943)
    at org.apache.spark.api.python.SerDeUtil$AutoBatchedPickler.foreach(SerDeUtil.scala:82)
    at org.apache.spark.api.python.PythonRDD$.writeIteratorToStream(PythonRDD.scala:442)
    at org.apache.spark.api.python.PythonRunner$$anon$2.writeIteratorToStream(PythonRunner.scala:871)
    at org.apache.spark.api.python.BasePythonRunner$WriterThread.$anonfun$run$1(PythonRunner.scala:573)
    at org.apache.spark.api.python.BasePythonRunner$WriterThread$$Lambda$2008/2134044540.apply(Unknown Source)
    at org.apache.spark.util.Utils$.logUncaughtExceptions(Utils.scala:2275)
    at org.apache.spark.api.python.BasePythonRunner$WriterThread.run(PythonRunner.scala:365)
22/09/07 10:25:52 ERROR SparkUncaughtExceptionHandler: Uncaught exception in thread Thread[stdout writer for /databricks/python/bin/python,5,main]
java.lang.OutOfMemoryError: Java heap space
    at org.apache.spark.sql.catalyst.expressions.UnsafeRow.getBinary(UnsafeRow.java:416)
    at org.apache.spark.sql.catalyst.expressions.SpecializedGettersReader.read(SpecializedGettersReader.java:75)
    at org.apache.spark.sql.catalyst.expressions.UnsafeRow.get(UnsafeRow.java:333)
    at org.apache.spark.sql.execution.python.EvaluatePython$.toJava(EvaluatePython.scala:58)
    at org.apache.spark.sql.execution.python.PythonForeachWriter.$anonfun$inputByteIterator$1(PythonForeachWriter.scala:43)
    at org.apache.spark.sql.execution.python.PythonForeachWriter$$Lambda$1830/1643360976.apply(Unknown Source)
    at scala.collection.Iterator$$anon$10.next(Iterator.scala:461)
    at org.apache.spark.api.python.SerDeUtil$AutoBatchedPickler.next(SerDeUtil.scala:92)
    at org.apache.spark.api.python.SerDeUtil$AutoBatchedPickler.next(SerDeUtil.scala:82)
    at scala.collection.Iterator.foreach(Iterator.scala:943)
    at scala.collection.Iterator.foreach$(Iterator.scala:943)
    at org.apache.spark.api.python.SerDeUtil$AutoBatchedPickler.foreach(SerDeUtil.scala:82)
    at org.apache.spark.api.python.PythonRDD$.writeIteratorToStream(PythonRDD.scala:442)
    at org.apache.spark.api.python.PythonRunner$$anon$2.writeIteratorToStream(PythonRunner.scala:871)
    at org.apache.spark.api.python.BasePythonRunner$WriterThread.$anonfun$run$1(PythonRunner.scala:573)
    at org.apache.spark.api.python.BasePythonRunner$WriterThread$$Lambda$2008/2134044540.apply(Unknown Source)
    at org.apache.spark.util.Utils$.logUncaughtExceptions(Utils.scala:2275)
    at org.apache.spark.api.python.BasePythonRunner$WriterThread.run(PythonRunner.scala:365)

集群配置

  • Driver节点:1个,14GB内存、4核CPU
  • Worker节点:2个,每个14GB内存、4核CPU

核心代码

cloudfile_options = {
    "cloudFiles.subscriptionId":subscription_ID,
    "cloudFiles.connectionString": queue_SAS_connection_string,
    "cloudFiles.format": "BinaryFile", 
    "cloudFiles.tenantId":tenant_ID,
    "cloudFiles.clientId":client_ID,
    "cloudFiles.clientSecret":client_secret,
    "cloudFiles.useNotifications" :"true"
}

def copy(row):
    source = row['path']
    destination = "somewhere"
    shutil.copy(source,destination)

spark.readStream.format("cloudFiles")
                        .options(**cloudfile_options)
                        .load(storage_input_path)              
                        .writeStream
                        .foreach(copy)
                        .option("checkpointLocation", checkpoint_location)
                        .trigger(once=True)
                        .start()

解决思路

1. 避免加载文件内容到Spark Row(根本解决)

使用BinaryFile格式时,Spark会将整个文件内容读取到DataFrame的content字段,大文件直接耗尽JVM堆内存。你的代码仅需文件路径,完全不需要加载内容,调整方案:

  • 移除cloudFiles.format": "BinaryFile"配置,改用cloudFiles.format": "text"(或其他轻量格式),并通过select("path")只获取文件路径字段
  • 若需过滤特定文件,可添加cloudFiles.filePattern参数指定匹配规则

修改后的读取代码示例:

spark.readStream.format("cloudFiles")
                .options(**cloudfile_options)
                .load(storage_input_path)
                .select("path")  # 仅保留路径字段,不加载文件内容
                .writeStream
                .foreach(copy)
                .option("checkpointLocation", checkpoint_location)
                .trigger(once=True)
                .start()

2. 替换为Databricks原生文件复制工具

用dbutils.fs.cp替代shutil.copy,该工具适配DBFS与Azure存储,无需将文件内容加载到内存,性能更优:

def copy(row):
    source = row['path']
    destination = "somewhere"
    dbutils.fs.cp(source, destination, recurse=False)

3. 调整Spark内存参数(临时缓解)

若必须保留BinaryFile格式,可调整Executor堆内存配置:

  • 在集群设置中,将spark.executor.memory设为10GB(节点总内存14GB,预留部分给系统进程)
  • 设置spark.executor.memoryOverhead为2GB以上,增加堆外内存配额
  • 调整spark.driver.memory至8GB,避免Driver端OOM

4. 限制单批次处理文件数量

通过maxFilesPerTrigger参数控制每次触发处理的文件数,避免同时加载过多大文件:

spark.readStream.format("cloudFiles")
                .options(**cloudfile_options, maxFilesPerTrigger=5)  # 按需调整数量
                .load(storage_input_path)
                .select("path")
                .writeStream
                .foreach(copy)
                .option("checkpointLocation", checkpoint_location)
                .trigger(once=True)
                .start()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.20 03:55:33