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

从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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 16:59:52