暂无图片
暂无图片
暂无图片
暂无图片
暂无图片

Spark-RDD算子-2

Coding On Road 2018-01-02
197

说明:

  目前发表的都为稿件,正式格式示例:




以下是依然稿件:

1、输入分区与输出分区合并累加型

 

1、union算子

功能:合并,但不去重复。

 

scala> declare #1 RDD 定义RDD1

scala> val rdd1 = sc.makeRDD(1 to 7,2);

rdd1: org.apache.spark.rdd.RDD[Int] = ParallelCollectionRDD[38] at makeRDD at <console>:24

//定义RDD2

scala> val rdd2 = sc.makeRDD(6 to 10,2);

rdd2: org.apache.spark.rdd.RDD[Int] = ParallelCollectionRDD[39] at makeRDD at <console>:24

//合并两个RDD,返回一个新的RDD

scala> val rdd3=rdd1.union(rdd2);

rdd3: org.apache.spark.rdd.RDD[Int] = UnionRDD[40] at union at <console>:28

//输出新的RDD的分区,结果为4个,即前两个分区的分区的和

scala> rdd3.partitions.length;

res58: Int = 4

//查看每一个分区的数据

scala> watch every partition data

scala> rdd3.mapPartitionsWithIndex((idx,p)=>{

     |    //declare a string to store data

     |    val str="Index is:"+idx+" data is:["+p.mkString(",")+"]";

     |    //must return iterator

     |    List(str).iterator;

     | }).collect();

//以下显示的每一个分区的数据

res59: Array[String] = Array(Index is:0 data is:[1,2,3], Index is:1 data is:[4,5,6,7], Index is:2 data is:[6,7], Index is:3 data is:[8,9,10])

 

图示:

可以看出,就是简单的累加。

如果两个RDD的分区不一样,也没有关系:

scala> val rdd4 = sc.makeRDD(4 to 9,3);

rdd4: org.apache.spark.rdd.RDD[Int] = ParallelCollectionRDD[42] at makeRDD at <console>:24

 

scala> val rdd5 = rdd1.union(rdd4);

rdd5: org.apache.spark.rdd.RDD[Int] = UnionRDD[43] at union at <console>:28

 

scala> rdd5.partitions.length;

res60: Int = 5

 

 

 

2、cartesian算子-笛卡尔集

功能:两个集合所有元素的积

源码:

def cartesian[U : ClassTag](other: RDD[U]): RDD[(T, U)]

     返回两个RDD的笛卡尔集:

示例1:都只有一个分区

分区中的数据一一进行组合。

scala> val rdd1 = sc.makeRDD(1 to 3,1);

rdd1: org.apache.spark.rdd.RDD[Int] = ParallelCollectionRDD[47] at makeRDD at <console>:24

 

scala> val rdd2 = sc.makeRDD(2 to 4,1);

rdd2: org.apache.spark.rdd.RDD[Int] = ParallelCollectionRDD[48] at makeRDD at <console>:24

 

scala> val rdd3 = rdd1.cartesian(rdd2);

rdd3: org.apache.spark.rdd.RDD[(Int, Int)] = CartesianRDD[49] at cartesian at <console>:28

 

scala> rdd3.collect

res63: Array[(Int, Int)] = Array((1,2), (1,3), (1,4), (2,2), (2,3), (2,4), (3,2), (3,3), (3,4))

 

图示:

示例2:都有两个分区

结果分区数量等于源所有分区数量的乘积。

scala> val rdd1 = sc.makeRDD(1 to 4,2); 定义第一个RDD

rdd1: org.apache.spark.rdd.RDD[Int] = ParallelCollectionRDD[50] at makeRDD at <console>:24

 

scala> rdd1.glom().collect; 显示每一个分区中的数据

res64: Array[Array[Int]] = Array(Array(1, 2), Array(3, 4))

 

scala> val rdd2 = sc.makeRDD(3 to 6,2); 定义第二个RDD

rdd2: org.apache.spark.rdd.RDD[Int] = ParallelCollectionRDD[52] at makeRDD at <console>:24

 

scala> rdd2.glom().collect;//显示每一个分区中的数据

res65: Array[Array[Int]] = Array(Array(3, 4), Array(5, 6))

 

scala> var rdd3 = rdd1.cartesian(rdd2); 执行笛卡尔集运算

rdd3: org.apache.spark.rdd.RDD[(Int, Int)] = CartesianRDD[54] at cartesian at <console>:28

 

scala> rdd3.collect; 输出结果

res66: Array[(Int, Int)] = Array((1,3), (1,4), (2,3), (2,4), (1,5), (1,6), (2,5), (2,6), (3,3), (3,4), (4,3), (4,4), (3,5), (3,6), (4,5), (4,6))

 

图示:

操作两个RDD的还有:

intersection - 两个集合的交集,将在后面讲到。

subtract - 用于获取在前RDD中出现,但在后RDD中没有的元素。将在后面讲到。

 

3、输入与输出多对多类型

1、groupBy

功能:分组

会触发shuffle操作。

源码:

def groupBy[K](f: T => K)

(implicit kt: ClassTag[K]): RDD[(K, scala.Iterable[T])]

说明:

Return an RDD of grouped items. Each group consists of a key and a sequence of elements mapping to that key. The ordering of elements within each group is not guaranteed, and may even differ each time the resulting RDD is evaluated.

以下是groupBy的重载:

groupBy会根据给定的key值做key,并用相同的key值做value:Interator

 

如:

  Jack,Jack,Mary,Rose元素,在GroupBy以后的输出为:以下是伪代码:

  (Jack,Iterator(Jack,Jack)),(Mary,Iterator(Mary)),(Rose,Iterator(Rose))

图解:

示例1:使用group实现单词统计

首先声明一个字符串的集合:

scala> val rdd:RDD[String] = sc.makeRDD(Seq("Jack","Mary","Rose","Jack"));

rdd: org.apache.spark.rdd.RDD[String] = ParallelCollectionRDD[12] at makeRDD at <console>:27

//查看里面的数据

scala> rdd.collect;

res10: Array[String] = Array(Jack, Mary, Rose, Jack)

//进行groupBy

scala> val rdd2=rdd.groupBy(str=>str);

rdd2: org.apache.spark.rdd.RDD[(String, Iterable[String])] = ShuffledRDD[14] at groupBy at <console>:29

//输出结果,正是以给定的值当key,注意最后一个元素是两个Jack

scala> rdd2.collect();

res11: Array[(String, Iterable[String])] = Array((Mary,CompactBuffer(Mary)), (Rose,CompactBuffer(Rose)), (Jack,CompactBuffer(Jack, Jack)))

//现在就可以进行快速的字符统计了:

scala> rdd2.collect().foreach(kv=>{

     |    println(kv._1+"\t"+kv._2.size);

     | });

Mary    1

Rose    1

Jack    2

 

以下是使用groupBy统计文件中单词的个数,注意.size的使用:

scala> val rdd = sc.textFile("file:///D:/a/a.txt");

scala> rdd.flatMap(_.split("\\s+")).groupBy(str=>str).collect().foreach(kv=>{

     |   println(kv._1+"\t"+kv._2.size);

     | });

Hello   4

Alex    1

Mary    1

Rose    1

Jack    1

 

Scala示例代码:

package cn.wang
import cn.wang.utils.ScUtils
import org.apache.spark.SparkContext
import org.apache.spark.rdd.RDD
object Demo04_GroupBy {
  def main(args: Array[String]): Unit = {
    val sc: SparkContext = ScUtils.sparkContext();
    val rdd: RDD[String] = sc.makeRDD(Seq("Jack", "Mary", "Jack", "Rose", "Jack"));
    val rdd2 = rdd.groupBy(str => str);
    rdd2.collect().foreach(kv => {
      println("统计结果:" + kv._1 + "\t" + kv._2.size);
    });
    sc.stop();
  }
}

 

也可以将数据保存到文件中,注意下例代码中saveAsTextFile的用法:

package cn.wang
import cn.wang.utils.ScUtils
import org.apache.spark.SparkContext
import org.apache.spark.rdd.RDD
object Demo04_GroupBy {
  def main(args: Array[String]): Unit = {
    val sc: SparkContext = ScUtils.sparkContext();
    val rdd: RDD[String] = sc.makeRDD(Seq("Jack", "Mary", "Jack", "Rose", "Jack"));
    val rdd2 = rdd.groupBy(str => str);
    //Array[Sting]类型的数据,不能保存到文件中,只要RDD[String]才有saveAsTextFile方法
    val rdd3:RDD[String] = rdd2.map(kv=>{
        val str = kv._1+"\t"+kv._2.size;
        str;
    });
    rdd3.saveAsTextFile("file:///D:/a/a1");
    sc.stop();
  }
}

 

2groupByKey

功能:根据key值进行分组。用于处理KeyValue类型的数据

源代码:

org.apache.spark.rdd.PairRDDFunctions

def groupByKey(): RDD[(K, scala.Iterable[V])]

属性[Key,Value]形式的方法,只能接收key,value对型的RDD

 

示例:统计单词个数

先进行map再进行groupByKey

package cn.wang
import cn.wang.utils.ScUtils
import org.apache.spark.SparkContext
import org.apache.spark.rdd.RDD
object Demo05_GroupByKey {
  def main(args: Array[String]): Unit = {
    val sc: SparkContext = ScUtils.sparkContext();
    var rdd: RDD[String] = sc.makeRDD(Seq("Jack", "Jack", "Alex", "Mary"));
    val rdd2: RDD[(String, Iterable[Int])] = rdd.map((_, 1)).groupByKey();

    //注意这儿size的使用,此时的value值是一个Iterable可遍历对象

    rdd2.map(kv => (kv._1, kv._2.size)).saveAsTextFile("file:///D:/a/a2");
    sc.stop();
  }
}

结果:

Alex1

Jack2

Mary1

 

 

命令行示例:

//声明字符串数组:

scala> val rdd = sc.makeRDD(Seq("Jack","Mary","Jack","Jack","Alex"));

rdd: org.apache.spark.rdd.RDD[String] = ParallelCollectionRDD[5] at makeRDD at <console>:24

//进行groupBy示例

scala> rdd.map((_,1)).groupByKey().collect;

res6: Array[(String, Iterable[Int])] = Array((Alex,CompactBuffer(1)), (Mary,CompactBuffer(1)), (Jack,CompactBuffer(1, 1, 1)))

//结果统计,由于reduceByKey接收类型应该是Int所以,中间再用map处理一下

scala> rdd.map((_,1)).groupByKey().map(kv=>(kv._1,kv._2.size)).reduceByKey((v1,v2)=>v1+v2).collect;

res8: Array[(String, Int)] = Array((Alex,1), (Mary,1), (Jack,3))

 

图示:

(伪代码)

3recudeByKey

源代码:

org.apache.spark.rdd.PairRDDFunctions

def reduceByKey(partitioner: Partitioner,

                func: (V, V) => V): RDD[(K, V)]

属于KeyValue RDD

见前面的示例。


文章转载自Coding On Road,如果涉嫌侵权,请发送邮件至:contact@modb.pro进行举报,并提供相关证据,一经查实,墨天轮将立刻删除相关内容。

评论