Spark DataFrame操作长时间运行无结果问题排查与优化咨询
问题原因分析与优化方案
从你提供的代码和运行日志来看,这个问题的核心矛盾出在循环遍历championIds并反复执行union操作的逻辑上,结合日志显示作业完成却无结果返回的现象,具体原因拆解如下:
核心问题根源
- Lineage(依赖链)爆炸:每次循环对
withBansDf执行union,都会不断累加DataFrame的依赖关系。当championIds数量较多时,整个数据集的依赖链会变得异常冗长,Driver在解析执行计划或组装最终结果时,可能因内存不足或计算超时陷入无响应——虽然Spark UI显示所有任务完成,但Driver仍在后台处理复杂的依赖逻辑,无法返回结果。 - 重复计算与资源浪费:每次循环都要对原始
dataframe执行where过滤、groupBy分组,相当于重复扫描和处理原始数据。大量重复的shuffle和计算不仅拖慢运行时间,还会让Driver端积累过多中间元数据,进一步加剧内存压力。 - 结果集组装瓶颈:经过数十上百次
union后,withBansDf的分区数会持续累加,Driver收集结果时需要合并大量小分区数据,容易出现内存溢出或无响应的情况。
针对性优化方案
我们可以彻底重构代码逻辑,抛弃循环+union的模式,改用批量处理的方式解决问题:
1. 批量展开所有被ban的英雄
先把原始数据中的team1_bans和team2_bans数组展开,让每行对应一个被ban的英雄,一次性覆盖所有英雄的ban数据:
import org.apache.spark.sql.functions._ // 展开两队的ban英雄列表,统一格式 val bannedChampionsDf = dataframe .select( $"patch_game_version", explode($"team1_bans").alias("champion_id") ) .union( dataframe .select( $"patch_game_version", explode($"team2_bans").alias("champion_id") ) )
2. 一次性计算所有英雄的ban次数与ban率
基于展开后的DataFrame,直接分组统计每个英雄在每个版本的ban次数,再关联gamesByPatchDf计算ban率:
// 统计每个英雄每个版本的ban总次数 val banCountDf = bannedChampionsDf .groupBy($"patch_game_version", $"champion_id") .count() .withColumnRenamed("count", "banned_count") // 关联版本游戏总数,计算ban率 val banRateDf = banCountDf .join(gamesByPatchDf, Seq("patch_game_version")) .withColumn("ban_rate", $"banned_count" / $"games_count") // 关联inProgress得到最终结果 val withBansDf = inProgress.join(broadcast(banRateDf), Seq("champion_id", "patch_game_version"))
3. 辅助配置优化
配合代码重构,调整以下Spark配置提升稳定性:
- 增加Driver内存:比如设置
--driver-memory 8g(根据实际资源调整),避免Driver处理大结果集或复杂逻辑时内存不足。 - 优化shuffle分区数:设置
spark.sql.shuffle.partitions 100(默认200,可根据数据量灵活调整),减少过多小分区带来的开销。 - 开启自适应执行计划:设置
spark.sql.adaptive.enabled true,让Spark自动调整执行计划和分区数,适配数据分布。
重构后的逻辑彻底避免了循环union带来的依赖链问题,同时大幅减少了重复计算,运行效率和稳定性都会显著提升,处理大量championIds时也能正常返回结果。
内容的提问来源于stack exchange,提问作者johnnydonna
相关产品推荐
相关产品推荐

