【发布时间】:2017-10-09 10:28:29
【问题描述】:
这是我正在做的:
val rddkv = sc.parallelize(List(("k1",1),("k2",2),("k1",2),("k3",5),("k3",1)))
//rddkv.collect
//Array[(String, Int)] = Array((k1,1), (k2,2), (k1,2), (k3,5), (k3,1))
rddkv.repartitionAndSortWithinPartitions(new org.apache.spark.RangePartitioner(3,rddkv)).mapPartitionsWithIndex( (i,iter_p) => iter_p.map(x=>" index="+i+" value="+x)).collect
//Array[String] = Array(" index=0 value=(k1,1)", " index=0 value=(k1,2)", " index=1 value=(k2,2)", " index=1 value=(k3,5)", " index=1 value=(k3,1)")
请注意,分区内的值未排序。这是为什么?我错过了什么?
【问题讨论】:
标签: scala apache-spark