【发布时间】:2020-08-12 00:50:15
【问题描述】:
我有一个 Spark 作业,它连接 2 个数据集,执行一些转换并减少数据以提供输出。 现在的输入大小非常小(每个 200MB 数据集),但是在加入之后,正如您在 DAG 中看到的那样,该作业被卡住并且永远不会继续进行第 4 阶段。我试着等了几个小时,它给了 OOM 并显示了第 4 阶段的失败任务。
- 为什么 spark 在 stage-3(join 阶段)之后不显示 stage-4(数据转换阶段)为活动状态?是不是陷入了第 3 阶段和第 4 阶段的洗牌?
- 如何提高 Spark 作业的性能?我尝试增加随机分区,结果仍然相同。
职位代码:
joinedDataset.groupBy("group_field")
.agg(collect_set("name").as("names")).select("names").as[List[String]]
.rdd. //converting to rdd since I need to use reduceByKey
.flatMap(entry => generatePairs(entry)) // this line generates pairs of words out of input text, so data size increases here
.map(pair => ((pair._1, pair._2), 1))
.reduceByKey(_+_)
.sortBy(entry => entry._2, ascending = false)
.coalesce(1)
仅供参考我的集群有 3 个 16 核和 100GB RAM 的工作节点,3 个 16 核的执行器(为简单起见,与机器的比例为 1:1)和 64GB 内存分配。
更新: 事实证明,我的工作中产生的数据非常庞大。我做了一些优化(战略性地减少输入数据并从处理中删除了一些重复的字符串),现在工作在 3 小时内完成。第 4 阶段的输入为 200MB,输出本身为 200GB。它正确地使用了并行性,但它在洗牌时很糟糕。我在这项工作中的 shuffle 溢出是 1825 GB(内存)和 181 GB(磁盘)。有人可以帮助我减少洗牌溢出和工作持续时间吗?谢谢。
【问题讨论】:
-
为什么要合并 (1)?
-
只是因为我想在 1 个文件中输出,没有其他原因。
-
对我来说很好用,但火花版本不同。
-
我可以知道你尝试了什么吗?您的应用程序在 stage4 之前挂起吗?
-
跑了一些东西 - 小卷 - 在数据块上。我不相信答案。你能用更小的体积作为测试运行吗
标签: scala apache-spark apache-spark-sql