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

Scala Spark使用自定义Case Class调用distinct.count报堆内存溢出

问题根因
  • 全量分组触发大Key内存溢出:groupBy(element => element.download_id)会将所有日志按用户ID分组,若存在高活跃用户对应的日志量级极大,单Key对应的所有数据会被加载到同一分区内存,超出堆内存上限就会触发OOM;你还对分组后的RDD做了缓存,进一步挤占执行内存。
  • 无效计算开销过高:你仅需要统计ghtorrent-22单个用户的数据,全量构建倒排索引属于不必要的操作,浪费大量内存和计算资源。
优化方案

逻辑重构(最优方案)

无需做全量分组和关联,直接在解析日志阶段就过滤目标用户,再做去重统计,从根源上避免大Key问题:

import org.apache.spark.rdd.RDD   
import java.text.SimpleDateFormat 
import java.util.Date
case class LogLine(debug_level: String, timestamp: Date, download_id: String,retrieval_stage: String, rest: String)
val regex = """([^\s]+), ([^\s]+)\+00:00, ([^\s]+) -- ([^\s]+): (.*$)""".r
// 声明目标用户,解析阶段直接过滤
val targetUserId = "ghtorrent-22"
val rdd = sc.textFile("file.txt")
  .flatMap ( x => x match {
      case regex(debug_level,dateTime,downloadId,retrievalStage,rest) if downloadId == targetUserId =>
        val df = new SimpleDateFormat("yyyy-MM-dd:HH:mm:ss")
        Some(LogLine(debug_level, df.parse(dateTime.replace("T", ":")), downloadId, retrievalStage, rest))
      case _ => None
  })
// 直接统计去重仓库数
val result = rdd.map(_.rest).distinct().count()

倒排索引逻辑兼容方案

如果业务需要复用全量倒排索引,可以做以下优化避免OOM:

  • 增大shuffle分区数,拆分大Key所在分区的数据量:将groupBy改为rdd.groupBy(_.download_id, numPartitions = 400),分区数可根据总数据量调整为默认值200的2-4倍
  • 切换为Kryo序列化减少对象内存占用,启动作业时添加配置:
    --conf spark.serializer=org.apache.spark.serializer.KryoSerializer
    --conf spark.kryo.registrationRequired=false
    

内存参数调优

  • 本地运行时,启动spark-shell添加参数--driver-memory 8G,提升Driver内存上限
  • 集群运行时,调整Executor配置:--executor-memory 8G --conf spark.executor.memoryOverhead=2048

内容的提问来源于stack exchange,提问作者Singularity

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.05 11:18:03