【问题标题】:What does df.repartition with no column arguments partition on?没有列参数的 df.repartition 分区是什么?
【发布时间】:2018-11-29 00:04:07
【问题描述】:

在 PySpark 中,重新分区模块有一个可选的列参数,它当然会通过该键重新分区您的数据帧。

我的问题是 - 当没有密钥时,Spark 如何重新分区?我无法进一步挖掘源代码以找到它通过 Spark 本身的位置。

def repartition(self, numPartitions, *cols):
    """
    Returns a new :class:`DataFrame` partitioned by the given partitioning expressions. The
    resulting DataFrame is hash partitioned.

    :param numPartitions:
        can be an int to specify the target number of partitions or a Column.
        If it is a Column, it will be used as the first partitioning column. If not specified,
        the default number of partitions is used.

    .. versionchanged:: 1.6
       Added optional arguments to specify the partitioning columns. Also made numPartitions
       optional if partitioning columns are specified.

    >>> df.repartition(10).rdd.getNumPartitions()
    10
    >>> data = df.union(df).repartition("age")
    >>> data.show()
    +---+-----+
    |age| name|
    +---+-----+
    |  5|  Bob|
    |  5|  Bob|
    |  2|Alice|
    |  2|Alice|
    +---+-----+
    >>> data = data.repartition(7, "age")
    >>> data.show()
    +---+-----+
    |age| name|
    +---+-----+
    |  2|Alice|
    |  5|  Bob|
    |  2|Alice|
    |  5|  Bob|
    +---+-----+
    >>> data.rdd.getNumPartitions()
    7
    """
    if isinstance(numPartitions, int):
        if len(cols) == 0:
            return DataFrame(self._jdf.repartition(numPartitions), self.sql_ctx)
        else:
            return DataFrame(
                self._jdf.repartition(numPartitions, self._jcols(*cols)), self.sql_ctx)
    elif isinstance(numPartitions, (basestring, Column)):
        cols = (numPartitions, ) + cols
        return DataFrame(self._jdf.repartition(self._jcols(*cols)), self.sql_ctx)
    else:
        raise TypeError("numPartitions should be an int or Column")

例如:调用这些行完全没问题,但我不知道它实际上在做什么。它是整行的哈希吗?也许是数据框中的第一列?

df_2 = df_1\
       .where(sf.col('some_column') == 1)\
       .repartition(32)\
       .alias('df_2')

【问题讨论】:

  • 我的回答是否帮助您了解了没有密钥的重新分区是如何工作的?

标签: python apache-spark pyspark pyspark-sql


【解决方案1】:

默认情况下,如果没有指定分区器,则分区不是基于数据的特性,而是以随机和均匀的方式分布在节点之间。

df.repartition 背后的重新分区算法会进行完整的数据洗牌,并在分区之间平均分配数据。为了减少洗牌,最好使用df.coalesce

这里有一些很好的解释如何用DataFrame重新分区 https://medium.com/@mrpowers/managing-spark-partitions-with-coalesce-and-repartition-4050c57ad5c4

【讨论】:

  • 所以它只使用行号?如果我能在源代码中得到它的引用,那就太好了。
  • 据我所知,它不使用您数据集中的任何信息,没有哈希键,它只是以均匀分布的方式重新分配数据(每个分区具有相同的大小)这是有道理的,甚至其他框架(如 apache kafka)不需要密钥来分区数据。如果没有提供密钥,Apache Kafka 分区数据默认使用 Round Robin
  • @Stefan Repcek,有什么办法可以根据数据框的总大小对数据进行分区吗?即假设我需要每个分区 128m ,如果我有 1GB ,我必须重新分区 1024 ,如果我有 5GB ,我必须重新分区 5120 ....所以应该动态计算重新分区数。
  • 链接失效。我认为这是正确的:medium.com/@mrpowers/…
猜你喜欢
  • 2017-03-17
  • 2023-01-19
  • 1970-01-01
  • 2010-09-14
  • 2018-11-29
  • 2022-11-17
  • 2017-04-19
  • 1970-01-01
相关资源
最近更新 更多