运行PySpark任务时遭遇OOM:ARRAY SIZE EXCEEDS VM LIMIT问题求助
问题分析与解决方案
错误本质
java array size exceeds the VM limit 核心是JVM单个数组长度不能超过Integer.MAX_VALUE(约21亿),你的问题根源在于那条1GB的单条记录:
- 1GB纯文本大致对应10亿字符,用
regexp_replace替换\n为\r\n时,会先把整条记录加载为内存字符串数组,替换后长度增加约50%;加上Java中char占2字节,该字符串内存占用接近3GB,超出Spark默认字符串长度限制(1GB),同时触发JVM数组长度阈值; - 重分区到30个分区无法解决问题,因为自定义分隔符
1%%会让这条大记录完整落在一个分区内,不会被拆分。
可行解决方案
1. 调整Spark字符串长度限制
Spark默认限制字符串最大长度为1GB,需修改参数放宽限制:
// SparkSession初始化时配置 spark.conf.set("spark.sql.maxStringLength", "2147483647") // 设为Integer.MAX_VALUE
或提交任务时添加:
--conf spark.sql.maxStringLength=2147483647
2. 改用逐行流式替换,避免全量内存加载
不要一次性加载整条大记录,用mapPartitions结合Java IO逐行处理:
import java.io.{BufferedReader, StringReader, StringWriter} df.mapPartitions(iter => { iter.map(row => { val originalContent = row.getString(0) val reader = new BufferedReader(new StringReader(originalContent)) val writer = new StringWriter() var line: String = null while ({line = reader.readLine(); line != null}) { writer.write(line) writer.write("\r\n") // 直接替换换行符 } val replacedContent = writer.toString() // 按业务逻辑返回Row,比如保留report_id字段 Row(replacedContent, row.getAs[String]("report_id")) }) })
这种方式逐行读写,不会把1GB记录一次性加载到内存,大幅降低内存压力。
3. 拆分超大记录(业务允许时)
如果业务逻辑允许拆分大记录,可在读取阶段处理:
- 自定义输入格式,将1GB记录按固定行数或大小拆分为子记录,同时保留
report_id标识; - 处理完成后再按
report_id合并回完整记录。
4. 优化Executor内存配置
调整Executor内存分配比例,增加堆外内存预留:
--conf spark.executor.memoryOverhead=4096m # 堆外内存设为4G,默认是Executor内存的10% --conf spark.memory.fraction=0.8 # 提高堆内存中用于存储/执行的比例
内容的提问来源于stack exchange,提问作者Raj harini
相关产品推荐
相关产品推荐

