在HDP集群用PySpark处理1GB文件遇Java堆内存不足的解决方法
首先咱们得先捋清楚问题根源:你用sc.wholeTextFiles()把整个文件直接加载成单个字符串,这会让每个完整文件都塞进Executor的内存里——1GB的JSON文件转成Java字符串后,内存占用会比原文件大不少(毕竟对象本身有额外开销),直接超出了默认的Executor堆内存限制,所以才会触发OutOfMemoryError。
下面给你几个针对性的解决方案,按优先级排序:
1. 放弃wholeTextFiles,用Spark原生JSON读取(最推荐)
Spark的JSON数据源本身就支持读取目录下的多个文件,而且会自动将文件拆分到多个分区,完全不需要把整个文件加载到内存里。把你的代码改成:
json_file = sqlContext.read.json('/user/admin/emp/*')
这个方法既简单又高效,除非你的JSON文件是每个文件只有一个超大JSON对象(不是每行一条JSON记录),否则这个方案直接就能解决问题。
2. 调整Executor和Driver的堆内存(如果必须用wholeTextFiles)
如果你的业务场景确实需要用wholeTextFiles(比如单文件对应单个完整JSON对象),那得给Spark分配足够的内存:
- 临时生效:启动pyspark时直接指定参数,根据你的集群资源调整数值,比如:
注意:Executor内存至少要比你最大的单个文件大1.5-2倍,因为字符串对象的内存开销不小。pyspark --executor-memory 4g --driver-memory 2g - 永久生效:在HDP的Ambari控制台里找到Spark服务的配置,修改
spark.executor.memory和spark.driver.memory的默认值,保存后重启Spark服务即可。
3. 增加分区数分散内存压力
如果你的目录下有多个文件,但默认分区数太少导致每个分区加载了太多文件,也会引发内存问题。可以在wholeTextFiles里指定最小分区数,让数据分散到更多分区:
json_file = sqlContext.read.json(sc.wholeTextFiles('/user/admin/emp/*', minPartitions=20).values())
分区数可以根据你的集群节点数和核数来调整,一般建议每个分区的数据量控制在100-200MB左右。
4. 优化Spark内存分配比例
如果调整了堆内存还是有问题,可以微调Spark的内存管理参数,让执行阶段有更多内存可用:
在启动pyspark时加上:
pyspark --conf spark.memory.fraction=0.8 --conf spark.memory.storageFraction=0.2
spark.memory.fraction设置了用于执行和存储的内存占堆内存的比例(默认0.6),spark.memory.storageFraction设置了存储内存占这个比例的部分(默认0.5),调整后能给执行内存留出更多空间。
内容的提问来源于stack exchange,提问作者kumar

