【问题标题】:Spark RDD partition by key in exclusive waySpark RDD 以独占方式按键分区
【发布时间】:2019-04-22 07:39:13
【问题描述】:

我想按键对 RDD 进行分区,并让每个分区仅包含单个键的值。例如,如果我有 100 个不同的键值和我 repartition(102),则 RDD 应该有 2 个空分区和 100 个分区,每个分区都包含一个键值。

我尝试使用 groupByKey(k).repartition(102) 但这并不能保证每个分区中键的排他性,因为我看到一些分区包含多个单个键的值并且超过 2 个空。

标准 API 中有没有办法做到这一点?

【问题讨论】:

    标签: apache-spark pyspark rdd


    【解决方案1】:

    要使用 partitionBy() RDD 必须由元组(对)对象组成。让我们看一个例子:

    假设我有一个包含以下数据的输入文件:

    OrderId|OrderItem|OrderDate|OrderPrice|ItemQuantity
    1|Gas|2018-01-17|1895|1
    1|Air Conditioners|2018-01-28|19000|3
    1|Television|2018-01-11|45000|2
    2|Gas|2018-01-17|1895|1
    2|Air Conditioners|2017-01-28|19000|3
    2|Gas|2016-01-17|2300|1
    1|Bottle|2018-03-24|45|10
    1|Cooking oil|2018-04-22|100|3
    3|Inverter|2015-11-02|29000|1
    3|Gas|2014-01-09|2300|1
    3|Television|2018-01-17|45000|2
    4|Gas|2018-01-17|2300|1
    4|Television$$|2018-01-17|45000|2
    5|Medicine|2016-03-14|23.50|8
    5|Cough Syrup|2016-01-28|190|1
    5|Ice Cream|2014-09-23|300|7
    5|Pasta|2015-06-30|65|2
    
    PATH_TO_FILE="file:///u/vikrant/OrderInputFile"
    

    将文件读入RDD并跳过标题

    RDD = sc.textFile(PATH_TO_FILE)
    header=RDD.first();
    newRDD = RDD.filter(lambda x:x != header)
    

    现在让我们将 RDD 重新分区为“5”个分区

    partitionRDD = newRDD.repartition(5)
    

    让我们看看数据在这“5”个分区中是如何分布的

    print("Partitions structure: {}".format(partitionRDD.glom().collect()))
    

    这里可以看到数据写入了两个分区,其中三个是空的,而且分布不均匀。

    Partitions structure: [[], 
    [u'1|Gas|2018-01-17|1895|1', u'1|Air Conditioners|2018-01-28|19000|3', u'1|Television|2018-01-11|45000|2', u'2|Gas|2018-01-17|1895|1', u'2|Air Conditioners|2017-01-28|19000|3', u'2|Gas|2016-01-17|2300|1', u'1|Bottle|2018-03-24|45|10', u'1|Cooking oil|2018-04-22|100|3', u'3|Inverter|2015-11-02|29000|1', u'3|Gas|2014-01-09|2300|1'], 
    [u'3|Television|2018-01-17|45000|2', u'4|Gas|2018-01-17|2300|1', u'4|Television$$|2018-01-17|45000|2', u'5|Medicine|2016-03-14|23.50|8', u'5|Cough Syrup|2016-01-28|190|1', u'5|Ice Cream|2014-09-23|300|7', u'5|Pasta|2015-06-30|65|2'], 
    [], []]
    

    我们需要创建一对 RDD,以使 RDD 数据均匀分布在多个分区中。 让我们创建一个pair RDD并将其分解为键值对。

    pairRDD = newRDD.map(lambda x :(x[0],x[1:]))
    

    现在让我们将此 rdd 重新分区为“5”分区,并使用第 [0] 位置的键将数据均匀分布到分区中。

    newpairRDD = pairRDD.partitionBy(5,lambda k: int(k[0]))
    

    现在我们可以看到数据正在根据匹配的键值对均匀分布。

    print("Partitions structure: {}".format(newpairRDD.glom().collect()))
    Partitions structure: [
    [(u'5', u'|Medicine|2016-03-14|23.50|8'), 
    (u'5', u'|Cough Syrup|2016-01-28|190|1'), 
    (u'5', u'|Ice Cream|2014-09-23|300|7'), 
    (u'5', u'|Pasta|2015-06-30|65|2')],
    
    [(u'1', u'|Gas|2018-01-17|1895|1'), 
    (u'1', u'|Air Conditioners|2018-01-28|19000|3'), 
    (u'1', u'|Television|2018-01-11|45000|2'), 
    (u'1', u'|Bottle|2018-03-24|45|10'), 
    (u'1', u'|Cooking oil|2018-04-22|100|3')], 
    
    [(u'2', u'|Gas|2018-01-17|1895|1'), 
    (u'2', u'|Air Conditioners|2017-01-28|19000|3'), 
    (u'2', u'|Gas|2016-01-17|2300|1')], 
    
    [(u'3', u'|Inverter|2015-11-02|29000|1'), 
    (u'3', u'|Gas|2014-01-09|2300|1'), 
    (u'3', u'|Television|2018-01-17|45000|2')], 
    
    [(u'4', u'|Gas|2018-01-17|2300|1'), 
    (u'4', u'|Television$$|2018-01-17|45000|2')]
    ]
    

    您可以在下面验证每个分区中的记录数。

    from pyspark.sql.functions import desc
    from pyspark.sql.functions import spark_partition_id
    
    partitionSizes = newpairRDD.glom().map(len).collect();
    
    [4, 5, 3, 3, 2]
    

    请注意,当您创建键值对的 pair RDD 时,您的键应该是 int 类型,否则您会收到错误。

    希望这会有所帮助!

    【讨论】:

    • 嘿维克兰特! partitionBy()repartition() 有什么区别?你不能互换使用它们吗?而不是在newpairRDD = pairRDD.partitionBy(5,lambda k: int(k[0])) 中使用partitionBy(),难道你没有在这里使用repartition() 吗?能详细说说两者的区别吗?
    • @cph_sto.. 是的,你可以。您可以在下面提到的链接中获得更多信息。 stackoverflow.com/questions/40416357/… & stackoverflow.com/questions/33831561/…
    • 维克兰特,在过去的几个月里,我已经阅读了您的许多问题/答案,它们非常有帮助。请允许我,尽管迟到了。
    • 这个link 很好地解决了我最初想到的问题。非常感谢您向我推荐这些链接。
    • @cph_sto.. 谢谢
    【解决方案2】:

    对于 RDD,您是否尝试过使用 partitionBy 对 RDD 进行密钥分区,例如 this question?如果需要,您可以将分区数指定为删除空分区的键数。

    在数据集 API 中,您可以使用 repartitionColumn 作为参数来按该列中的值进行分区(但请注意,这使用 spark.sql.shuffle.partitions 的值作为分区数,因此您'会得到更多的空分区)。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2017-06-15
      • 1970-01-01
      • 1970-01-01
      • 2021-02-24
      • 1970-01-01
      • 2019-06-16
      • 2017-02-17
      相关资源
      最近更新 更多