【问题标题】:Why the Spark's repartition didn't balance data into partitions?为什么 Spark 的重新分区没有将数据平衡到分区中?
【发布时间】:2019-09-12 11:23:51
【问题描述】:
>>> rdd = sc.parallelize(range(10), 2)
>>> rdd.glom().collect()
[[0, 1, 2, 3, 4], [5, 6, 7, 8, 9]]
>>> rdd.repartition(3).glom().collect()
[[], [0, 1, 2, 3, 4], [5, 6, 7, 8, 9]]
>>>

第一个分区是空的?为什么?非常感谢您告诉我原因。

【问题讨论】:

  • 您是担心第一个分区为空还是某个分区为空?

标签: apache-spark pyspark rdd


【解决方案1】:

值得注意的是,由于 Spark 就是大规模运行,因此不太可能担心这种情况。您可以获得的最接近的是倾斜数据。 range 将给出与使用散列的重新分区不同的初始分区。对批量大小的评论也是有效的,但在实践中不太相关。

【讨论】:

    【解决方案2】:

    这可以通过查看重新分区功能的工作原理来解释。 这样做的原因是调用df.repartition(COL, numPartitions=k) 将使用基于哈希的分区创建一个带有k 分区的数据帧。 Pyspark 将遍历每一行并应用以下function 来确定当前行中的元素将在哪里结束:

    partition_the_row_belongs_to = hash(COL) % k
    

    在这种情况下,k 用于将行映射到由 k 个分区组成的空间。如您所见,哈希函数有时会发生冲突。有时有些分区是空的,而其他分区的元素太多。这可能是因为散列图结论,也可能是因为散列函数。无论哪种方式,您所看到的原因是 repartition 按照您的要求创建了 3 个分区,它不会向您保证任何关于平衡分区或让所有分区都非空的事情。如果您想更好地控制生成的分区的外观,请查看partitionby。

    另请参阅:this question 和 this question。

    我希望这会有所帮助。

    【讨论】:

      【解决方案3】:

      这是因为 Spark 不会随机播放单个元素,而是随机播放数据块 - with minimum batch size equal to 10。

      因此,如果您的元素少于每个分区的元素,Spark 将不会分离分区的内容。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2021-01-26
        • 2018-11-14
        • 1970-01-01
        • 2012-09-28
        • 1970-01-01
        • 2019-11-13
        • 1970-01-01
        相关资源
        最近更新 更多