本地IntelliJ运行Scala Spark脚本时遭遇内存溢出问题
解决Spark本地模式合并Excel文件时的OOM问题
问题分析
你遇到的java.lang.OutOfMemoryError: Java heap space,核心原因有两点:
- 本地模式(
local[*])下,Spark的Executor与Driver共享同一个JVM进程,spark.executor.memory配置完全无效,所有内存开销都由Driver承担。 - 循环读取单个Excel再手动
union的方式,会在Driver端堆积多个DataFrame实例,加上Excel解析本身的内存开销,容易触发内存溢出。
具体解决方案
1. 修正SparkSession配置(针对本地模式)
删掉无效的spark.executor.memory,重点优化Driver相关内存参数:
val session = SparkSession.builder() .config("spark.driver.bindAddress", "127.0.0.1") .config("spark.driver.memory", "12g") // 充分利用32G机器资源,直接给到12G .config("spark.driver.maxResultSize", "4g") // 增大结果数据内存上限,避免show/collect时溢出 .config("spark.sql.shuffle.partitions", "8") // 本地模式调小shuffle分区,减少内存开销 .config("spark.memory.offHeap.enabled", true) .config("spark.memory.offHeap.size", "4g") .master("local[*]") .appName("etl") .getOrCreate()
2. 优化Excel读取逻辑(避免手动union)
不要循环读取单个文件,直接让Spark读取整个目录下的所有Excel文件,Spark会自动并行处理,避免Driver端堆积数据:
// 替换原有的文件遍历和read_excel函数 val mdf = session.read.excel(header = true) .schema(dataSchema) .load("data/*.xlsx") // 直接读取目录下所有xlsx文件
3. 确保IntelliJ的VM参数配置正确
在IntelliJ的Run Configuration中,修改VM Options:
-Xms8g -Xmx16g
该参数用于设置Driver进程的JVM初始堆和最大堆,需与Spark配置的spark.driver.memory匹配或更大(Spark的driver内存基于JVM堆分配)。
修改后的完整代码
import com.crealytics.spark.excel._ import org.apache.spark.sql.types.StructField import org.apache.spark.sql.{SparkSession, types} object sparkJob extends App { val session = SparkSession.builder() .config("spark.driver.bindAddress", "127.0.0.1") .config("spark.driver.memory", "12g") .config("spark.driver.maxResultSize", "4g") .config("spark.sql.shuffle.partitions", "8") .config("spark.memory.offHeap.enabled", true) .config("spark.memory.offHeap.size", "4g") .master("local[*]") .appName("etl") .getOrCreate() val dataSchema = types.StructType(Array( StructField("Delivery Date", types.StringType, nullable = false), StructField("Delivery Hour", types.IntegerType, nullable = false), StructField("Delivery Interval", types.IntegerType, nullable = false), StructField("Repeated Hour Flag", types.StringType, nullable = false), StructField("Settlement Point Name", types.StringType, nullable = false), StructField("Settlement Point Type", types.StringType, nullable = false), StructField("Settlement Point Price", types.DecimalType(10, 0), nullable = false) )) // 直接读取目录下所有Excel文件,自动合并 val mdf = session.read.excel(header = true) .schema(dataSchema) .load("data/*.xlsx") mdf.show(5) session.stop() // 记得关闭SparkSession,释放资源 }
额外提示
- 提前检查所有Excel文件的列名、数据类型是否与定义的schema完全匹配,格式不匹配会导致解析时额外的内存消耗。
- 本地模式下尽量避免使用
collect()等将全量数据拉到Driver的操作,show(5)只取前5条没问题,但后续处理全量数据时,尽量用Spark的分布式操作。
内容的提问来源于stack exchange,提问作者ndl4145
相关产品推荐
相关产品推荐

