【问题标题】:Spark: increase number of partitions without causing a shuffle?Spark:增加分区数量而不引起洗牌?
【发布时间】:2015-01-18 07:02:42
【问题描述】:

当减少分区数量时,可以使用coalesce,这非常棒,因为它不会导致随机播放并且似乎可以立即工作(不需要额外的工作阶段)。

有时我想反其道而行之,但repartition 会引发洗牌。我想几个月前我实际上是通过使用CoalescedRDDbalanceSlack = 1.0 来完成这项工作的——所以会发生什么情况是它会拆分一个分区,以便生成的分区位于同一个节点上(如此小的网络IO)。

这种功能在 Hadoop 中是自动的,只需调整拆分大小。除非减少分区的数量,否则它在 Spark 中似乎不会以这种方式工作。我认为解决方案可能是编写一个自定义分区器以及一个自定义 RDD,我们在其中定义 getPreferredLocations ......但我认为这是一件如此简单和常见的事情,肯定必须有一种直接的方法吗?

尝试过的事情:

.set("spark.default.parallelism", partitions) 在我的SparkConf 上,并且在阅读镶木地板的情况下我尝试过sqlContext.sql("set spark.sql.shuffle.partitions= ...,它在 1.0.0 上会导致错误并且不是我想要的,我希望分区号改变所有类型的工作,而不仅仅是洗牌。

【问题讨论】:

  • 运气好能找到解决方案吗?

标签: scala apache-spark


【解决方案1】:

注意这个空间

https://issues.apache.org/jira/browse/SPARK-5997

这种非常简单明显的功能最终会实现——我猜他们在完成Datasets 中所有不必要的功能之后。

【讨论】:

    【解决方案2】:

    我不完全明白你的意思。你的意思是你现在有 5 个分区,但是在下一次操作之后你希望数据分布到 10 个?因为有 10 个,但仍然使用 5 个没有多大意义……将数据发送到新分区的过程必须在某个时候发生。

    在执行coalesce 时,您可以摆脱未使用的分区,例如:如果您最初有 100 个,但在 reduceByKey 之后您得到了 10 个(因为那里只有 10 个键),您可以设置 coalesce

    如果您希望进程以另一种方式进行,您可以强制进行某种分区:

    [RDD].partitionBy(new HashPartitioner(100))
    

    我不确定这就是你要找的东西,但希望如此。

    【讨论】:

    • 每个分区都有一个位置,即一个节点,假设我有5个分区和5个节点。如果我将repartition 或您的代码调用到 10 个分区,这将打乱数据 - 即 5 个节点中的每一个的数据都可能通过网络传递到其他节点。我想要的是,Spark 只是将每个分区分成 2 个而不移动任何数据——这就是 Hadoop 在调整拆分设置时发生的情况。
    • 我不确定你能不能做到。我猜你需要某种.forEachNode 函数。但我从未见过这样的事情。而且我不确定它是否可以轻松实施。分区器每次都必须为同一个对象返回相同的分区。默认情况下,Spark 使用 HashPartitioner,它执行 hashCode modulo number_of_partitions。如果您只是将数据拆分为两个新分区,那么它们肯定会最终出现在它们的位置上。这就是为什么洗牌是必要的。也许如果你有自己的分区器,它可以增加分区的数量而不用在网络上洗牌。
    【解决方案3】:

    如您所知,pyspark 使用某种“懒惰”的运行方式。它只会在需要执行某些操作时进行计算(例如“df.count()”或“df.show()”。所以您可以做的是定义这些操作之间的随机分区。

    你可以写:

    sparkSession.sqlContext().sql("set spark.sql.shuffle.partitions=100")
    # you spark code here with some transformation and at least one action
    df = df.withColumn("sum", sum(df.A).over(your_window_function))
    df.count() # your action
    
    df = df.filter(df.B <10)
    df = df.count()   
    
    sparkSession.sqlContext().sql("set spark.sql.shuffle.partitions=10")
    # you reduce the number of partition because you know you will have a lot 
    # less data
    df = df.withColumn("max", max(df.A).over(your_other_window_function))
    df.count() # your action
    

    【讨论】:

    • spark.sql.shuffle.partitions只会对连接、聚合和排序等混排操作产生影响......但对过滤没有影响
    猜你喜欢
    • 2018-01-29
    • 1970-01-01
    • 1970-01-01
    • 2016-04-29
    • 2020-03-20
    • 1970-01-01
    • 2019-07-01
    • 2015-04-08
    • 2021-06-14
    相关资源
    最近更新 更多