使用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
相关产品推荐
相关产品推荐

