spark RDD算子(九)之基本的Action操作 first, take, collect, count, countByValue, reduce, aggregate, fold,top
·
章节目录
一、first
返回第一个元素
Scala版本
scala> val rdd = sc.parallelize(List(1,2,3,3,4))
rdd: org.apache.spark.rdd.RDD[Int] = ParallelCollectionRDD[4] at parallelize at <console>:24
scala> rdd.first()
res3: Int = 1

Java版本
// first
JavaRDD<Integer> firstRDD = sc.parallelize(Arrays.asList(1, 2, 3, 4));
Integer first = firstRDD.first();
System.out.println(first);

二、take
rdd.take(n) 返回前n个元素
Scala版本
scala> val rdd = sc.parallelize(List(1,2,3,3,4))
rdd: org.apache.spark.rdd.RDD[Int] = ParallelCollectionRDD[5] at parallelize at <console>:24
scala> rdd.take(4)
res4: Array[Int] = Array(1, 2, 3, 3)

Java版本
// take
JavaRDD<Integer> takeRDD = sc.parallelize(Arrays.asList(1, 2, 3, 3, 4, 5));
List<Integer> take = takeRDD.take(3);
System.out.println(take);

三、collect
rdd.collect() 返回RDD中的所有元素
Scala版本
scala> val rdd = sc.parallelize(List(1,2,3,3,4))
rdd: org.apache.spark.rdd.RDD[Int] = ParallelCollectionRDD[6] at parallelize at <console>:24
scala> rdd.collect()
res5: Array[Int] = Array(1, 2, 3, 3, 4)

Java版本
// collect
JavaRDD<Integer> collectRDD = sc.parallelize(Arrays.asList(1, 2, 3, 4, 5));
List<Integer> collect = collectRDD.collect();
System.out.println(collect);

四、count
rdd.count() 返回RDD的元素个数
Scala版本
scala> val rdd = sc.parallelize(List(1,2,3,3,4))
rdd: org.apache.spark.rdd.RDD[Int] = ParallelCollectionRDD[7] at parallelize at <console>:24
scala> rdd.count()
res7: Long = 5

Java版本
// count
JavaRDD<Integer> countRDD = sc.parallelize(Arrays.asList(1, 2, 3, 4, 5, 6));
long count = countRDD.count();
System.out.println(count);

五、countByValue
各元素在RDD中出现的次数,返回{(key1,次数),(key2,次数),…(key3,次数)}
Scala版本
scala> val rdd = sc.parallelize(List(1,2,3,3,4))
rdd: org.apache.spark.rdd.RDD[Int] = ParallelCollectionRDD[9] at parallelize at <console>:24
scala> rdd.countByValue()
res9: scala.collection.Map[Int,Long] = Map(4 -> 1, 1 -> 1, 3 -> 2, 2 -> 1)

Java版本
// countByValue
JavaRDD<Integer> countByValueRDD = sc.parallelize(Arrays.asList(1, 2, 3, 3, 4, 4, 5));
Map<Integer, Long> integerLongMap = countByValueRDD.countByValue();
System.out.println(integerLongMap);

六、reduce
rdd.reduce(func)
并行整合RDD中所有的数据,类似于Scala中的集合的reduce
Scala版本
scala> val rdd = sc.parallelize(List(1,2,3,3,4))
rdd: org.apache.spark.rdd.RDD[Int] = ParallelCollectionRDD[13] at parallelize at <console>:24
scala> rdd.reduce((x,y)=>x+y)
res10: Int = 13

Java版本
// reduce
JavaRDD<Integer> rdd = sc.parallelize(Arrays.asList(1, 2, 3, 4, 5));
Integer reduceRDD = rdd.reduce(new Function2<Integer, Integer, Integer>() {
@Override
public Integer call(Integer integer, Integer integer2) throws Exception {
return integer + integer2;
}
});
System.out.println(reduceRDD);

七、aggregate
aggregate函数将每个分区里面的元素进行聚合,然后用combine函数将每个分区的结果和初始值(zeroValue)进行combine操作。这个函数最终返回的类型不需要和RDD中元素类型一致。
注意:
1.每个分区开始聚合第一个元素都是zeroValue
2.分区之间的聚合,zeroValue也参与运算
Scala版本
//分区数为1时
val conf = new SparkConf().setMaster("local[1]").setAppName("aggregate")
val sc = new SparkContext(conf)
val rdd = sc.parallelize(List(1,2,3))
def seqop(x:Int,y:Int)={
println("seqop",x,y)
x+y
}
def combine(x:Int,y:Int)={
println("combine",x,y)
x+y
}
val result = rdd.aggregate(1)(seqop,combine)
println(result)

//分区数为2时
val conf = new SparkConf().setMaster("local[2]").setAppName("aggregate")
val sc = new SparkContext(conf)
val rdd = sc.parallelize(List(1,2,3))
def seqop(x:Int,y:Int)={
println("seqop",x,y)
x+y
}
def combine(x:Int,y:Int)={
println("combine",x,y)
x+y
}
val result = rdd.aggregate(1)(seqop,combine)
println(result)

八、fold
rdd.fold(num)(func) 一般不用这函数
和reduce()一样,但是提供了初始值num,每个元素计算时,先要和这个初始值进折叠,注意:这里会按照每个分区进行fold,然后分区之间再次进行fold
Scala版本
scala> val rdd = sc.parallelize(List(1,2,3),2)
rdd: org.apache.spark.rdd.RDD[Int] = ParallelCollectionRDD[1] at parallelize at <console>24
scala> rdd.fold(1)((x,y)=>{println(x,y);x+y})
(1,1)
(1,2)
(3,3)
(1,2)
(3,6)
res4: Int = 9

Java版本
// fold
JavaRDD<Integer> rdd1 = sc.parallelize(Arrays.asList(1, 2, 3, 3), 2);
Integer foldRDD = rdd1.fold(1, new Function2<Integer, Integer, Integer>() {
@Override
public Integer call(Integer v1, Integer v2) throws Exception {
return v1 + v2;
}
});
System.out.println(foldRDD);

九、top
rdd.top(n) 按照降序或者指定的排序规则排序,并返回前n各元素
Scala版本
scala> val rdd = sc.parallelize(List(1,2,3,3))
rdd: org.apache.spark.rdd.RDD[Int] = ParallelCollectionRDD[5] at parallelize at <console>:24
scala> rdd.top(3)
res5: Array[Int] = Array(3, 3, 2)

Java版本
// top
JavaRDD<Integer> topRDD = sc.parallelize(Arrays.asList(1, 2, 3, 3));
List<Integer> top = topRDD.top(2);
System.out.println(top);

十、takeOrdered
rdd.takeOrdered(n)
对RDD元素进行升序排序,取出前n个元素并返回,也可以自定义比较器(这里不介绍),类似于top的相反的方法
Scala版本
scala> val rdd = sc.parallelize(List(1,2,3,3,4))
rdd: org.apache.spark.rdd.RDD[Int] = ParallelCollectionRDD[7] at parallelize at <console>:24
scala> rdd.takeOrdered(3)
res7: Array[Int] = Array(1, 2, 3)

Java版本
// takeOrdered
JavaRDD<Integer> takeOrderedRDD = sc.parallelize(Arrays.asList(1, 2, 3, 3),2);
List<Integer> takeOrdered = takeOrderedRDD.takeOrdered(2);
System.out.println(takeOrdered);

十一、foreach
对RDD中的每个元素使用给定的函数
Scala版本
val rdd = sc.parallelize(List(1,2,3,3,4))
rdd.foreach(print(_))

Java版本
// foreach
JavaRDD<Integer> foreachdRDD = sc.parallelize(Arrays.asList(1, 2, 3, 3,4));
foreachdRDD.foreach(new VoidFunction<Integer>() {
@Override
public void call(Integer integer) throws Exception {
System.out.print(integer);
}
});

更多推荐
所有评论(0)