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

以下是依然稿件:
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();
}
}
2:groupByKey
功能:根据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))
图示:
(伪代码)

3:recudeByKey
源代码:
org.apache.spark.rdd.PairRDDFunctions
def reduceByKey(partitioner: Partitioner,
func: (V, V) => V): RDD[(K, V)]
属于KeyValue RDD。
见前面的示例。




