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

运行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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.05 09:55:16