【问题标题】:Why is the number of partitions after groupBy 200? Why is this 200 not some other number?为什么groupBy后面的分区数是200?为什么这 200 不是其他数字?
【发布时间】:2017-01-17 03:52:45
【问题描述】:

这是 Spark 2.2.0-SNAPSHOT。

为什么下例中groupBy变换后的分区数为200?

scala> spark.range(5).groupByKey(_ % 5).count.rdd.getNumPartitions
res0: Int = 200

200 有什么特别之处?为什么不是像1024 这样的其他号码?

有人告诉我 Why does groupByKey operation have always 200 tasks? 专门询问 groupByKey,但问题是选择 200 作为默认值背后的“谜团”,而不是为什么默认有 200 个分区。

【问题讨论】:

标签: apache-spark


【解决方案1】:

这是由 spark.sql.shuffle.partitions 设置的

一般来说,每当您执行 spark sql 聚合或连接以打乱数据时,这就是结果分区的数量。

对于您的整个操作来说,它是恒定的(即,不可能为一个转换更改它,然后再为另一个转换更改它)。

更多信息请参见http://spark.apache.org/docs/latest/sql-programming-guide.html#other-configuration-options

【讨论】:

  • 我不认为有什么特别的原因。他们似乎有 200 个东西,但与某种优化无关。我认为需要制定一个基准,以便从中获得一些东西
  • 一般来说,spark 默认值旨在在本地计算机(1 个节点)中开箱即用。在调整集群时,它们应该配置得更好。
  • spark.shuffle.sort.bypassMergeThreshold 也是 200。我不认为这是巧合。
  • 文档链接现在是spark.apache.org/docs/latest/…
猜你喜欢
  • 2017-06-03
  • 2012-11-13
  • 1970-01-01
  • 2018-12-11
  • 1970-01-01
  • 2010-10-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多