【问题标题】:Spark repartition is slow and shuffles too much dataSpark重新分区很慢并且洗牌太多数据
【发布时间】:2015-05-05 17:52:21
【问题描述】:

我的集群:

  • 5个数据节点
  • 每个数据节点有:8 个 CPU,45GB 内存

由于其他一些配置限制,我只能在每个数据节点上启动 5 个执行器。所以我做到了

spark-submit --num-executors 30 --executor-memory 2G ...

所以每个执行器使用 1 个核心。

我有两个数据集,每个大约 20 GB。在我的代码中,我做到了:

val rdd1 = sc...cache()
val rdd2 = sc...cache()

val x = rdd1.cartesian(rdd2).repartition(30) map ...

在 Spark UI 中,我看到 repartition 步骤耗时 30 多分钟,导致数据 shuffle 超过 150GB。

我认为这是不对的。但我不知道出了什么问题...

【问题讨论】:

  • 顺便说一句,你应该总是在重新分区后缓存,否则你最终会在每次点击时随机播放。

标签: apache-spark


【解决方案1】:

您真的是指“笛卡尔”吗?

您将 RDD1 中的每一行乘以 RDD2 中的每一行。因此,如果您的行是每行 1k,那么每个 RDD 大约有 20,000 行。笛卡尔积将返回一个包含 20,000 x 20,000 或 4 亿条记录的集合。请注意,现在每行的宽度将加倍——2k——所以你在 RDD3 中有 800 GB,而在 RDD1 和 RDD2 中只有 20 GB。

或许可以试试:

val x = rdd1.union(rdd2).repartition(30) map ...

甚至可能:

val x = rdd1.zip(rdd2).repartition(30) map ...

?

【讨论】:

  • 顺便说一下,如果您的记录大小更像我的世界中的记录大小——100 字节左右——每个文件的记录数将是其十倍,每个文件大约有 20 万条记录。在笛卡尔积之后,您将拥有 400 亿条记录,每条 200 字节,用于 8 TB 数据。
  • 是的,我的意思是 catesian。我知道它会变大,但没想到只是做分区就这么慢。
  • 我想说 30 分钟实际上是非常快的。您预期的速度有多快?为什么?
猜你喜欢
  • 2019-07-01
  • 2016-04-29
  • 2021-04-23
  • 1970-01-01
  • 2015-01-18
  • 1970-01-01
  • 2019-07-25
  • 1970-01-01
  • 2020-12-16
相关资源
最近更新 更多