Spark 2.0.2应用调优求助:解决OOM与Executor心跳超时异常
结合你的Spark 2.0.2环境、m4.2xlarge集群配置,以及当前遇到的两类异常,我从资源参数调优、代码逻辑优化两方面给你具体的解决方案:
一、先解决资源分配不合理的问题
你的集群节点是16vCPU、32GiB内存,当前的Executor参数设置存在资源过度分配、内存拆分不合理的问题,这是导致两个异常的核心原因之一:
修正Executor核数与数量
m4.2xlarge节点需要预留2-3核给系统进程,单节点可用CPU核数建议控制在13以内。你当前设置--num-executors 5 --executor-cores 5,总核数25远超单节点承载能力,会导致Executor跨节点分散、资源竞争激烈,进而引发GC延迟和心跳超时。
建议调整为:--num-executors 2 --executor-cores 6单节点跑2个Executor,每个分配6核,总核数12,既充分利用CPU资源,又给系统留足缓冲空间。
优化内存分配,解决Container内存超限
你当前--executor-memory 9g + memoryOverhead=1g的组合,堆内存占比过高,而Spark的Off-heap内存、依赖包加载、系统进程都会消耗Overhead内存,容易触发YARN的内存限制。建议调整为:--executor-memory 8g --conf spark.yarn.executor.memoryOverhead=2048降低堆内存占比,把Overhead提升到2G,每个Executor总内存仍为10G,但分配更合理,避免被YARN强制杀死。
缓解Executor心跳超时
心跳超时大多是Executor长时间GC导致无法响应YARN心跳,除了内存调整,还可以放宽心跳限制:--conf spark.yarn.executor.heartbeatInterval=10000 --conf spark.yarn.max.executor.failures=10把心跳间隔从默认1000ms延长到10000ms,同时允许更多次Executor失败重试,避免因临时GC超时直接终止任务。
二、优化代码逻辑,从根源减少内存压力
你的代码中allUsers.cartesian(allGears)是最大的内存杀手——这个操作会生成用户总数×装备总数的爆炸式RDD,缓存后直接占满Executor内存,必须重构:
替换笛卡尔积操作,避免全量数据膨胀
你最终需要的是用户未交互的装备,完全不用先做全量笛卡尔积再取差集,换一种更高效的方式:// 先把已购买/加购的装备按用户分组,缓存这个小体量的RDD val gearShoppingGrouped = gearShoppingUserToItem.groupByKey().persist(StorageLevel.MEMORY_AND_DISK_SER) // 把全量装备转成广播变量,每个Executor只加载一份,节省内存 val allGearsBroadcast = sc.broadcast(allGears.map(_.toInt).collect().toSet) // 生成有交互记录用户的未交互装备 val nonPurchasedWithInteraction = gearShoppingGrouped.flatMap{ case (userId, purchasedItems) => val purchasedSet = purchasedItems.toSet allGearsBroadcast.value.diff(purchasedSet).map(gearId => (userId, gearId)) } // 处理无任何交互记录的用户,生成他们的全量装备候选 val allUsersSet = allUsers.map(_.toInt).collect().toSet val noInteractionUsers = allUsersSet.diff(gearShoppingGrouped.keys.collect().toSet) val nonPurchasedNoInteraction = sc.parallelize(noInteractionUsers.toSeq).flatMap(userId => allGearsBroadcast.value.map(gearId => (userId, gearId)) ) // 合并最终的候选集 val finalNonPurchasedGears = nonPurchasedWithInteraction.union(nonPurchasedNoInteraction)这种方式用广播变量传递全量装备,避免了笛卡尔积的海量数据,内存占用会大幅降低。
优化缓存策略
- 只缓存分组后的
gearShoppingGrouped,这个RDD数据量远小于全量笛卡尔积结果; - 缓存时指定
MEMORY_AND_DISK_SER级别,内存不足时自动溢写到磁盘,避免OOM。
- 只缓存分组后的
及时清理无用数据
在执行model.predict前,确保所有不再使用的RDD都执行unpersist(),比如你之前的allUserItems.unpersist()操作要提前到该RDD不再使用时执行,不要等到最后。
三、其他辅助调优建议
- 启用Kryo序列化
替换默认的Java序列化,减少内存占用和序列化时间:--conf spark.serializer=org.apache.spark.serializer.KryoSerializer --conf spark.kryoserializer.buffer.max=512m - 使用G1垃圾收集器
减少GC停顿时间,避免心跳超时:--conf spark.executor.extraJavaOptions="-XX:+UseG1GC -XX:+PrintGCDetails -XX:+PrintGCTimeStamps" - 广播大模型
如果你的训练模型体积较大,用sc.broadcast(model)把模型广播给所有Executor,避免每个Executor重复加载模型浪费内存。
最后,整合所有调整后的spark-submit命令:
spark-submit --jars jedis-2.7.2.jar,commons-pool2-2.3.jar,spark-redis-0.3.2.jar,SparkHBase.jar,recommendcontentslib_2.11-1.0.jar \ --class org.digitaljuice.itemrecommender.RecommendGears \ --master yarn \ --driver-memory 2g \ --num-executors 2 \ --executor-memory 8g \ --executor-cores 6 \ --conf spark.yarn.executor.memoryOverhead=2048 \ --conf spark.yarn.executor.heartbeatInterval=10000 \ --conf spark.yarn.max.executor.failures=10 \ --conf spark.serializer=org.apache.spark.serializer.KryoSerializer \ --conf spark.kryoserializer.buffer.max=512m \ --conf spark.executor.extraJavaOptions="-XX:+UseG1GC -XX:+PrintGCDetails -XX:+PrintGCTimeStamps" \ recommendersystem_2.11-0.0.1.jar /work/output/gearpurchaserating/part-00000 /work/output/gearaddtocartrating/part-00000 /work/output/allGears/part-00000 /work/output/allAccounts/part-00000 /work/allaccounts/acc_toacc/part-m-00000 /work/Recommendations/ /work/TrainingModel
按照这个方案调整后,应该能顺利解决内存超限和心跳超时问题,完成预测任务。
内容的提问来源于stack exchange,提问作者Adit

