【问题标题】:How to repartition RDD in spark?如何在火花中重新分区RDD?
【发布时间】:2016-04-23 08:06:46
【问题描述】:

我有一些数据如下:

b   3
c   1
a   1
b   2
b   1
a   2

我想按第一列重新分区3个部分,并保存为文件,结果应该是这样的(不需要排序):

//file: part-00000
a   1
a   2

//file: part-00001
b   3
b   2
b   1

//file: part-00002
c   1

我尝试调用重新分区函数,但无法达到我的目的。

怎么做?非常感谢!

【问题讨论】:

  • Spark 支持 Hash 和 Range Partitioner。虽然 Range Partitioner 可能在一定程度上满足您的需求,但如果两者都不能满足您的要求,那么您需要编写自定义 Partitioner。但是,即使在考虑之前,您能否说出按照您定义的方式进行分区的原因/好处?
  • 感谢您的回复!我知道如何定义自定义分区器,但是在调用repartition function时不知道如何使用,也不知道在哪里可以正确使用自定义分区器。你能告诉我如何使用自定义分区器吗?我的自定义分区器是:class MyPartitioner(val partitions: Int = 1) extends Partitioner{ def numPartitions = partitions def getPartition(key: Any): Int = { val s = key.toString if(s.equals("a")) 0 else if(s.equals("b")) 1 % partitions else if(s.equals("c")) 2 % partitions else 0 } }
  • 请看看答案是否适合你。
  • @郭;你也可以回答这个问题:stackoverflow.com/questions/31610971/…

标签: apache-spark


【解决方案1】:

Sumit 回答的更多补充:
实现您的自定义org.apache.spark.Partitioner。 例如:

class AlphbetPartitioner extends Partitioner {

  override def numPartitions: Int = 26

  override def getPartition(key: Any): Int = {

    return key.asInstanceOf[scala.Char].asDigit % numPartitions
  }
}

PairRDDFunctions.partitionBy(partitioner: Partitioner) 的示例代码

val data = Array(('b', 3), ('c', 1), ('a', 1), ('b', 2), ('b', 1), ('a', 2))
val distData = sc.parallelize(data,1).map(u => (u._1, u._2)).partitionBy(new AlphbetPartitioner).map(u=>u._1+","+u._2+"\t")

【讨论】:

  • 我明白了!谢谢你的回答!
【解决方案2】:

您需要调用partitionBy-函数来使用您的自定义分区器对数据进行分区。我可以推荐阅读这本在线书籍的“数据分区”部分:https://www.safaribooksonline.com/library/view/learning-spark/9781449359034/ch04.html

【讨论】:

    【解决方案3】:

    自定义分区器只能与 RDD 一起使用以键入键/值,即PairRDDFunctions.partitionBy(partitioner: Partitioner)。更多信息请参考PairRDDFunctions API

    【讨论】:

    • 谢谢!答案对我很有用。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-07-06
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多