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
相关产品推荐
相关产品推荐

