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

本地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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 14:07:43