一、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);
    }
});

在这里插入图片描述

Logo

腾讯云面向开发者汇聚海量精品云计算使用和开发经验,营造开放的云计算技术生态圈。

更多推荐