从Azure Data Lake读取Parquet文件到Databricks时出现Java堆内存溢出错误
我们有一个Azure Databricks Notebook,用于读取Azure Data Lake中"RAW"文件夹下的Parquet文件。这些文件由Azure Event Hub流式写入,按年/月/日分区存储,单日文件夹下包含10000+个大小在20KB到6MB之间的小文件。
示例文件路径:RAW/ConsumptionData/Parquet/V1/CapturedYear=2023/CapturedMonth=10/CapturedDay=30/137341728_872c9ee8bbaf4c58b8ee7a401f459b41_1.parquet
此前运行正常,现在执行时崩溃,报错:
java.lang.OutOfMemoryError: Java heap space
核心代码
读取数据的核心代码:
# 读取数据并指定Schema consFromRAW = ( spark .read .option("header","True") .schema(schema) .parquet( "/mnt/datalakegen2/RAW/ConsumptionData/Parquet/V1") )
后续过滤逻辑(因OOM无法执行到这一步):
# 过滤所需的日期分区 consFromRAWFiltered = consFromRAW.filter(col("CapturedDate") >= firstRAWReadDate)
其中firstRAWReadDate通常为当前时间减去2天,CapturedDate由分区字段CapturedYear、CapturedMonth、CapturedDay拼接而成。
集群配置
- 驱动节点:Standard_DS3_v2
- 工作节点:Standard_DS3_v2
- 节点数量:2-16
- Databricks版本:10.4 LTS(包含Apache Spark 3.2.1、Scala 2.12)
完整错误栈
java.lang.OutOfMemoryError: Java heap space --------------------------------------------------------------------------- Py4JJavaError Traceback (most recent call last) <command-3957504232682417> in <module> 72 73 # Read the data with the defined schema ---> 74 consFromRAW = ( spark 75 .read 76 .option("header","True") /databricks/spark/python/pyspark/sql/readwriter.py in parquet(self, *paths, **options) 299 int96RebaseMode=int96RebaseMode) 300 ---> 301 return self._df(self._jreader.parquet(_to_seq(self._spark._sc, paths))) 302 303 def text(self, paths, wholetext=False, lineSep=None, pathGlobFilter=None, /databricks/spark/python/lib/py4j-0.10.9.1-src.zip/py4j/java_gateway.py in __call__(self, *args) 1302 1303 answer = self.gateway_client.send_command(command) -> 1304 return_value = get_return_value( 1305 answer, self.gateway_client, self.target_id, self.name) 1306 /databricks/spark/python/pyspark/sql/utils.py in deco(*a, **kw) 115 def deco(*a, **kw): 116 try: -> 117 return f(*a, **kw) 118 except py4j.protocol.Py4JJavaError as e: 119 converted = convert_exception(e.java_exception) /databricks/spark/python/lib/py4j-0.10.9.1-src.zip/py4j/protocol.py in get_return_value(answer, gateway_client, target_id, name) 324 value = OUTPUT_CONVERTER[type](answer[2:], gateway_client) 325 if answer[1] == REFERENCE_TYPE: -> 326 raise Py4JJavaError( 327 "An error occurred while calling {0}{1}{2}.\n". 328 format(target_id, ".", name), value) Py4JJavaError: An error occurred while calling o533.parquet. : java.lang.OutOfMemoryError: Java heap space at scala.collection.mutable.HashTable.resize(HashTable.scala:258) at scala.collection.mutable.HashTable.addEntry0(HashTable.scala:158) at scala.collection.mutable.HashTable.findOrAddEntry(HashTable.scala:170) at scala.collection.mutable.HashTable.findOrAddEntry$(HashTable.scala:167) at scala.collection.mutable.LinkedHashSet.findOrAddEntry(LinkedHashSet.scala:44) at scala.collection.mutable.LinkedHashSet.add(LinkedHashSet.scala:68) at scala.collection.mutable.LinkedHashSet.$plus$eq(LinkedHashSet.scala:63) at scala.collection.mutable.LinkedHashSet.$plus$eq(LinkedHashSet.scala:44) at scala.collection.generic.Growable.$anonfun$$plus$plus$eq$1(Growable.scala:62) at scala.collection.generic.Growable$$Lambda$10/936292831.apply(Unknown Source) at scala.collection.mutable.ResizableArray.foreach(ResizableArray.scala:62) at scala.collection.mutable.ResizableArray.foreach$(ResizableArray.scala:55) at scala.collection.mutable.ArrayBuffer.foreach(ArrayBuffer.scala:49) at scala.collection.generic.Growable.$plus$plus$eq(Growable.scala:62) at scala.collection.generic.Growable.$plus$plus$eq$(Growable.scala:53) at scala.collection.mutable.AbstractSet.$plus$plus$eq(Set.scala:50) at org.apache.spark.sql.execution.datasources.InMemoryFileIndex.$anonfun$listLeafFiles$2(InMemoryFileIndex.scala:144) at org.apache.spark.sql.execution.datasources.InMemoryFileIndex$$Lambda$3969/1659440957.apply(Unknown Source) at scala.collection.mutable.ResizableArray.foreach(ResizableArray.scala:62) at scala.collection.mutable.ResizableArray.foreach$(ResizableArray.scala:55) at scala.collection.mutable.ArrayBuffer.foreach(ArrayBuffer.scala:49) at org.apache.spark.sql.execution.datasources.InMemoryFileIndex.listLeafFiles(InMemoryFileIndex.scala:140) at org.apache.spark.sql.execution.datasources.InMemoryFileIndex.refresh0(InMemoryFileIndex.scala:102) at org.apache.spark.sql.execution.datasources.InMemoryFileIndex.<init>(InMemoryFileIndex.scala:74) at org.apache.spark.sql.execution.datasources.DataSource.createInMemoryFileIndex(DataSource.scala:623) at org.apache.spark.sql.execution.datasources.DataSource.resolveRelation(DataSource.scala:458) at org.apache.spark.sql.DataFrameReader.loadV1Source(DataFrameReader.scala:356) at org.apache.spark.sql.DataFrameReader.$anonfun$load$2(DataFrameReader.scala:323) at org.apache.spark.sql.DataFrameReader$$Lambda$3794/1655878321.apply(Unknown Source) at scala.Option.getOrElse(Option.scala:189) at org.apache.spark.sql.DataFrameReader.load(DataFrameReader.scala:323) at org.apache.spark.sql.DataFrameReader.parquet(DataFrameReader.scala:740)
一、代码层面优化
1. 直接读取目标分区路径,避免全量扫描
既然只需要最近2天的数据,直接构造对应日期的分区路径,减少驱动节点需要扫描的文件数量:
from datetime import datetime, timedelta # 计算目标日期范围 end_date = datetime.now() start_date = end_date - timedelta(days=2) # 生成需要读取的分区路径列表 paths = [] current_date = start_date while current_date <= end_date: year = current_date.year month = current_date.month day = current_date.day path = f"/mnt/datalakegen2/RAW/ConsumptionData/Parquet/V1/CapturedYear={year}/CapturedMonth={month}/CapturedDay={day}" paths.append(path) # 读取指定路径的Parquet文件(Parquet无需header=True,可直接去掉) consFromRAW = ( spark .read .schema(schema) .parquet(*paths) )
2. 利用原生分区列过滤,而非拼接字段
直接使用Spark的分区发现机制,通过CapturedYear、CapturedMonth、CapturedDay过滤,让Spark提前裁剪分区:
consFromRAWFiltered = consFromRAW.filter( (col("CapturedYear") >= start_date.year) & (col("CapturedMonth") >= start_date.month) & (col("CapturedDay") >= start_date.day) )
若日期跨月跨年,需调整条件确保覆盖所有目标日期
二、集群配置调整
1. 升级驱动节点规格
错误栈显示OOM发生在InMemoryFileIndex阶段,是驱动节点扫描文件元数据时内存不足导致的。Standard_DS3_v2仅14GB内存,建议升级到Standard_DS4_v2(28GB内存)或更高规格。
2. 调整Spark参数优化内存使用
在Notebook开头添加以下配置:
# 开启Parquet分区下推 spark.conf.set("spark.sql.parquet.filterPushdown", "true") # 设置分区文件合并阈值(128MB),减少小文件数量 spark.conf.set("spark.sql.files.maxPartitionBytes", "134217728") # 调整文件打开成本阈值,让Spark优先合并小文件 spark.conf.set("spark.sql.files.openCostInBytes", "10485760") # 给驱动节点分配更多内存(需节点规格支持) spark.conf.set("spark.driver.memory", "10g")
三、数据存储层面优化
1. 定期合并历史分区小文件
添加Databricks定时任务,定期合并历史分区的小文件,提升后续读取效率:
# 合并指定路径下的小文件 def compact_parquet(path): df = spark.read.parquet(path) # 根据数据量调整分区数,确保输出文件在100-200MB之间 df.repartition(10).write.mode("overwrite").parquet(path) # 调用函数合并指定日期的分区 compact_parquet("/mnt/datalakegen2/RAW/ConsumptionData/Parquet/V1/CapturedYear=2023/CapturedMonth=10/CapturedDay=29")
2. 调整Event Hub写入策略
修改流式写入逻辑,减少小文件生成:
- 提高文件滚动大小阈值(如设置为100MB再生成新文件)
- 延长文件滚动时间间隔(如每30分钟生成一个文件)
- 合并多个微批次写入,避免频繁生成小文件
内容的提问来源于stack exchange,提问作者JonJagd

