习题 11:自定义分区器
题目:实现自定义 Partitioner,按 key 值取模分区;生成 100 万个随机数,按自定义分区器分区。
// UseridPartitioner.scala
import org.apache.spark.Partitioner
class UseridPartitioner(nump: Int) extends Partitioner {
override def numPartitions: Int = nump
override def getPartition(key: Any): Int = {
return key.toString().toInt % nump
}
}
// TestUserPartition.scala
object TestUserPartition {
def main(args: Array[String]): Unit = {
val conf = new SparkConf().setAppName("test").setMaster("local[*]")
val sc = new SparkContext(conf)
val rdd = sc.parallelize(a)
val rdd2 = rdd.map((_, 1)).partitionBy(new UseridPartitioner(6))
rdd2.collect()
sc.stop()
}
}