【问题标题】:Spark write to CSV fails even after 8 hours即使 8 小时后,Spark 写入 CSV 也会失败
【发布时间】: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

执行者页面

Spark Web UI 页面用于返回随机播放错误的节点

节点的 Spark Web UI 页面返回 ec2 连接关闭错误

总体工作摘要页面

【问题讨论】:

  • 感觉有点像笛卡尔连接,不是吗?
  • 您必须在 ec2 中写入小文件。我猜你正试图一次写入整个数据。
  • 您在回复 Ram 时写道,sourcedest 中有 200 个分区,但没有打印 result 中的分区数 - 有多少?另外,这里的笛卡尔连接会产生什么样的通货膨胀? 10 倍? 100 倍? 1000 倍?
  • 我的错误 Tim - 结果也有 200 个分区。使用重新分区,正如我在更新中所指出的,我已将其增加到 66K 分区。至于通货膨胀,我该如何衡量?源和目标的对应行数相同,因此结果数据帧最终的行数相同,但各有 2 列

标签: apache-spark spark-dataframe


【解决方案1】:

我可以从 Spark UI 中看到它甚至从未写入 CSV 阶段,忙于坚持阶段,但没有坚持 功能它在 8 小时后仍然失败。有什么想法吗?

它是FetchFailedException,即未能获取随机块

由于您能够处理小文件,因此只有大数据会失败... 我强烈觉得分区不够。

首先是验证/打印source.rdd.getNumPartitions()。和destinations.rdd.getNumPartitions()。和result.rdd.getNumPartitions()

您需要在加载数据后重新分区,以便将数据(通过 shuffle)分区到集群中的其他节点。这将为您提供更快处理而不会失败所需的并行性

进一步,验证应用的其他配置... 像这样打印所有配置,根据需要将它们调整为正确的值。

sc.getConf.getAll

也看看

【讨论】:

  • 抱歉耽搁了,我会尽快尝试您的建议
  • 这是我得到的输出:scala> sources.rdd.getNumPartitions res25: Int = 200 scala> destinations.rdd.getNumPartitions res26: Int = 200
  • getAll 命令的输出:pastebin.com/T8aYuBVd ...我尝试了重新分区和合并功能,但由于某种原因它总是返回 200。知道为什么它不起作用吗?
【解决方案2】:

在加入之前重新分区源和目标,分区数使每个分区为 10MB - 128MB(尝试调整),无需将其设为 20000(恕我直言太多)。 然后通过这两列连接然后写入,不重新分区(即输出分区应该与连接前的重新分区相同)

如果你仍然有问题,在将两个数据帧转换为 rdd 后尝试做同样的事情(api 之间存在一些差异,尤其是在重新分区、键值 rdds 等方面)

【讨论】:

    猜你喜欢
    • 2018-08-02
    • 1970-01-01
    • 2016-03-30
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-06-02
    • 1970-01-01
    相关资源
    最近更新 更多