【发布时间】:2017-11-13 03:11:37
【问题描述】:
我有一个数据框,其中包含大约 200-600 GB 的数据,我正在读取、操作,然后使用弹性映射减少集群上的 spark shell (scala) 写入 csv。即使 8 小时后,Spark 写入 CSV 也会失败
这是我写 csv 的方式:
result.persist.coalesce(20000).write.option("delimiter",",").csv("s3://bucket-name/results")
结果变量是通过来自其他一些数据帧的列的混合创建的:
var result=sources.join(destinations, Seq("source_d","destination_d")).select("source_i","destination_i")
现在,我可以在大约 22 分钟内读取它所基于的 csv 数据。在同一个程序中,我还能够在 8 分钟内将另一个(较小的)数据帧写入 csv。但是,对于这个 result 数据帧,它需要 8 多个小时并且仍然失败......说其中一个连接已关闭。
我也在 13 x c4.8xlarge instances on ec2 上运行这项工作,每个内核有 36 个内核和 60 GB 内存,所以我认为我有能力写入 csv,尤其是在 8 小时之后。
许多阶段需要重试或任务失败,我无法弄清楚我做错了什么或为什么需要这么长时间。我可以从 Spark UI 中看到,它甚至从未进入写入 CSV 阶段并且忙于持久化阶段,但没有持久化功能,它在 8 小时后仍然失败。有任何想法吗?非常感谢您的帮助!
更新:
我已运行以下命令将 result 变量重新分区为 66K 分区:
val r2 = result.repartition(66000) #confirmed with numpartitions
r2.write.option("delimiter",",").csv("s3://s3-bucket/results")
但是,即使在几个小时之后,作业仍然失败。我还做错了什么?
注意,我正在通过 spark-shell yarn --driver-memory 50G 运行 spark shell
更新 2:
我已经尝试先用持久化运行写入:
r2.persist(StorageLevel.MEMORY_AND_DISK)
但我有很多阶段失败,返回一个,Job aborted due to stage failure: ShuffleMapStage 10 (persist at <console>:36) has failed the maximum allowable number of times: 4. Most recent failure reason: org.apache.spark.shuffle.MetadataFetchFailedException: Missing an output location for shuffle 3' 或说Connection from ip-172-31-48-180.ec2.internal/172.31.48.180:7337 closed
【问题讨论】:
-
感觉有点像笛卡尔连接,不是吗?
-
您必须在 ec2 中写入小文件。我猜你正试图一次写入整个数据。
-
您在回复 Ram 时写道,
source和dest中有 200 个分区,但没有打印result中的分区数 - 有多少?另外,这里的笛卡尔连接会产生什么样的通货膨胀? 10 倍? 100 倍? 1000 倍? -
我的错误 Tim - 结果也有 200 个分区。使用重新分区,正如我在更新中所指出的,我已将其增加到 66K 分区。至于通货膨胀,我该如何衡量?源和目标的对应行数相同,因此结果数据帧最终的行数相同,但各有 2 列
标签: apache-spark spark-dataframe