【郑重声明:本文为原创作品,版权所有,如转载请注明出处。
如作为商业宣传,本人将保留诉讼的权力。--王健】
同时欢迎技术爱好者提出问题,共同进步。本文示展示的程序,都是用Scala语言编写。

【郑重声明:本文为原创作品,版权所有,如转载请注明出处。
如作为商业宣传,本人将保留诉讼的权力。】
Spark RDD编程
1、RDD的概念
RDD(Resilient Distributed Datasets) ,弹性分布式数据集, 是分布式内存的一个抽象概念,RDD提供了一种高度受限的共享内存模型,即RDD是只读的记录分区的集合,只能通过在其他RDD执行确定的转换操作(如map、join和group by)而创建,然而这些限制使得实现容错的开销很低。对开发者而言,RDD可以看作是Spark的一个对象,它本身运行于内存中,如读文件是一个RDD,对文件计算是一个RDD,结果集也是一个RDD ,不同的分片、 数据之间的依赖 、key-value类型的map数据都可以看做RDD。
2、RDD算了分类
1:转换算子 - 也是惰性的。
如:sc.textFile(..);并不会真实的加载这个文件
rdd.filter(line=>line.contains(“Jack”)); 不会真实的执行过虑
转换操作,返回的都是RDD。
转换操作不会执行JOB。
2:行动算了
如:
rdd.count();
rdd.collect(); 即,只有真实的使用这些数据时,才会去执行过虑或是计数任务。
行动操作会执行JOB。
强调:
RDD是不可变的,转换操作,是生成一个新的RDD而不是修改原有的RDD。
3、算子的功能
通过转换算子,获取一个新的RDD。
通过Action算子,触发Spark提交作业。
通过Cache算子,将数据缓存到内存,具体的说,应该是缓存到内存的堆空间。
在运行转换中通过算子对RDD进行转换。算子是RDD中定义的函数,可以对RDD中的数据进行转换和操作。
1)输入:在Spark程序运行中,数据从外部数据空间(如分布式存储:textFile读取HDFS等,parallelize方法输入Scala集合或数据)输入Spark,数据进入Spark运行时数据空间,转化为Spark中的数据块,通过BlockManager进行管理。
2)运行:在Spark数据输入形成RDD后便可以通过变换算子,如fliter等,对数据进行作并将RDD转化为新的RDD,通过Action算子,触发Spark提交作业。 如果数据需要复用,可以通过Cache算子,将数据缓存到内存。
3)输出:程序运行结束数据会输出Spark运行时空间,存储到分布式存储中(如saveAsTextFile输出到HDFS),或Scala数据或集合中(collect输出到Scala集合,count返回Scala int型数据)。Spark的核心数据模型是RDD,但RDD是个抽象类,具体由各子类实现,如MappedRDD、 ShuffledRDD等子类。 Spark将常用的大数据操作都转化成为RDD的子类。
4、RDD算子
Spark算子是由Scala写成的,所以,看懂Scala语法很重要。在讲解以下算子时,本人会附上大量的Scala的源代码,帮助大家理解。
1、转换算子
转换算子,生成新的RDD。不触发Job。但有些转换算子如repartition会有shuffle操作,但依然不会执行job。只有行动算子,才会触发Job。
1、输入分区与输出分区一对一类型
即输出时的分区与原分区保持相同:
//设置为分7个分区如果文件正好7个字节的话,否则将会根据算法,重新计算分区//数量,见Hadoop的FileInputFormat的getSplits方法源代码(旧版),因到到作者写此文//章时,spark2.1.2依然使用的是Hadoop的mapred下的类,即旧版本的api。
scala> val rdd = sc.textFile("file:///D:/a/a.txt",7);
rdd: org.apache.spark.rdd.RDD[String] = file:///D:/a/a.txt MapPartitionsRDD[16] at textFile at <console>:24
scala> rdd.partitions.length; //显示为7个分区
res16: Int = 7
scala> val rdd2 = rdd.map(_.split("\\s+")); ///进行转换,不会触发Job作业
rdd2: org.apache.spark.rdd.RDD[Array[String]] = MapPartitionsRDD[17] at map at <console>:26
scala> rdd2.partitions.length; 转换以后还是7个分区
res18: Int = 7
1、map算子
是转换算子的一种。
将一个RDD中的每个数据项,通过map中的函数映射变为一个新的元素。返回的是一个数组。
输入分区与输出分区一对一,即:有多少个输入分区,就有多少个输出分区。
map算子的源代码说明:
见map算子的源代码:
def map[U: ClassTag](f: T => U): RDD[U] = withScope {
val cleanF = sc.clean(f)
new MapPartitionsRDD[U, T](this, (context, pid, iter) => iter.map(cleanF))
}
说明:
(f: T => U):中的f为接收的变量,T为这个变量的类型
如:”Jack Mary”.map(name:String=>...),其中name就是指的:Jack Mary,它显然是String类型。
(f: T => U) :中的U就是你需要返回的类型是什么,如果你这样写:
“Jack Mary”.map(name:String = > name.split(“ “));则当然返回的是数组类型。即:RDD[Array[String]]。
如果你这样写:
“Jack Mary”.map(name:String=>name.replace(“ “,””));即去除里面的空格。则这个map函数,就返回字符串数组即:RDD[String]。
(f: T => U): RDD[U] 这里面的RDD[U],即U正是你的返回类型。将返回一个包含U类型的新的RDD对象。
通过上面的说明,你应该可以明白,返回什么类型完全由你程序执行的结果来定。
读其他源码,也是类似的思路。慢慢的练习,就可以了。
示例1:返回Array[Array[Sting]]
有以下文本文件a.txt,内容如下:
HelloJack
HelloMary
HelloRose
HelloAlex
以下代码,将按行显示每一行map以后的结果。
scala> val rdd1 = sc.textFile("file:///D:/a/a.txt"); //按行读取文件中的数据,获取一个新RDD
rdd1: org.apache.spark.rdd.RDD[String] = file:///D:/a/a.txt MapPartitionsRDD[3]..at <console>:24
scala> val rdd2 =rdd1.map(line=>line.split("\\s+"));//使用map对第一个行进行分割
rdd2: org.apache.spark.rdd.RDD[Array[String]] = MapPartitionsRDD[4]..<console>:26
scala> rdd2.collect(); 显示数据,注意与flatMap的区别
res3: Array[Array[String]] = Array(Array(Hello, Jack), Array(Hello, Mary), Array(Hello, Rose), Array(Hello, Alex))
Scala代码:
package cn.wang
import cn.wang.utils.ScUtils
import org.apache.spark.SparkContext
import org.apache.spark.rdd.RDD
object Demo01_Map {
def main(args: Array[String]): Unit = {
//获取sc对象
val sc: SparkContext = ScUtils.sparkContext();
//读取文件
val rdd: RDD[String] = sc.textFile("file:///D:/a/a.txt");
//执行转换
val rdd2: RDD[Array[String]] = rdd.map(_.split("\\s+"));
//输出数据[Array(Hello,Jack),Array(Hello,Mary)...]
rdd2.collect().foreach(arr => {
//由于每一个都是一个数组array所以可以通过apply获取具体的值
println("这个数组中的数据是:" + arr.apply(0) + "," + arr.apply(1));
});
sc.stop();
}
}
输出的结果:
这个数组中的数据是:Hello,Jack
这个数组中的数据是:Hello,Mary
这个数组中的数据是:Hello,Rose
这个数组中的数据是:Hello,Alex
注意:map方法,返回的是一个(A,B)数组
示例2:返回Array[Array[Int]]
再如如下的示例:
package cn.wang
import cn.wang.utils.ScUtils
import org.apache.spark.SparkContext
import org.apache.spark.rdd.RDD
object Demo01_01_map {
def main(args: Array[String]): Unit = {
val sc: SparkContext = ScUtils.sparkContext();
val rdd: RDD[Int] = sc.parallelize(List(1, 2, 3));
//注意map算子的返回值是:RDD[(X,Y)]
val rdd2: RDD[(Int, Int)] = rdd.map(age => {
val aa: Int = age + 1;
(aa, 100);//注意100是一个固定值,只是为了示例没有实际功能
});
rdd2.foreach(println(_));
sc.stop();
}
}
/**
* 输出的结果:
* (2,100)
(3,100)
(4,100)
*/
示例3:返回Array[String]
package cn.wang
import cn.wang.utils.ScUtils
import org.apache.spark.SparkContext
import org.apache.spark.rdd.RDD
object Demo01_02_map {
def main(args: Array[String]): Unit = {
val sc: SparkContext = ScUtils.sparkContext();
val rdd: RDD[String] = sc.textFile("file:///D:/a/a.txt");
println("************以下map方法返回的是Array[String]*************");
val rdd2: RDD[Array[String]] =
rdd.map(str => {
val ss: Array[String] = str.split("\\s+");
//由于这儿做了一个Split所以,返回的是数组Array[String]
ss;
});
rdd2.collect().foreach(abc => {
println(abc.apply(0) + "," + abc.apply(1));
});
println("*************以下的map方法返回一String类型************************");
val rdd3: RDD[String] =
rdd.map(str => {
//这儿是将里面的空格去掉,即:Jack Mary => JackMary,所以,返回的是String类型
str.replaceAll("\\s+", "");
});
rdd3.collect().foreach(println(_));
sc.stop();
}
}
注意上面第二个map算子的使用,由于它返回的是str.replaceAll(..)即返回的是一个字符串,所以这个RDD的返回类型为:Array[String].
说明:
可见,Map方法,具体返回的类型则map算子最后最后的返回值来决定。
2:flatMap算子
flatMap属于Transformation算子,第一步和map一样,最后将所有的输出分区合并成一个。
示例1:显示每一个单词
命令行:
scala> val rdd = sc.textFile("file:///D:/a/a.txt");
rdd: org.apache.spark.rdd.RDD[String] = file:///D:/a/a.txt MapPartitionsRDD[8] at textFile at <console>:24
//使用flatMap组成一个Array(str,str,str.....)的结构
scala> val rdd2 = rdd.flatMap(line=>line.split("\\s+"));
rdd2: org.apache.spark.rdd.RDD[String] = MapPartitionsRDD[9] at flatMap at <console>:26
scala> rdd2.collect();
res11: Array[String] = Array(Hello, Jack, Hello, Mary, Hello, Rose, Hello, Alex)
//输出数据
scala> rdd2.collect().foreach(str=>{
| println("字符"+str);
| }
| );
字符串是Hello
字符串是Jack
字符串是Hello
字符串是Mary
字符串是Hello
字符串是Rose
字符串是Hello
字符串是Alex
Scala代码:
package cn.wang
import cn.wang.utils.ScUtils
import org.apache.spark.SparkContext
import org.apache.spark.rdd.RDD
object Demo02_flatMap {
def main(args: Array[String]): Unit = {
val sc: SparkContext = ScUtils.sparkContext();
val rdd: RDD[String] = sc.textFile("file:///D:/a/a.txt");
val rdd2: RDD[String] = rdd.flatMap(_.split("\\s+"));
rdd2.collect().foreach(str => {
println("值:" + str);
});
sc.stop();
}
}
示例2:单词统计
以下示例,用flatMap与map的组合实现单词的统计:
package cn.wang
import cn.wang.utils.ScUtils
import org.apache.spark.SparkContext
import org.apache.spark.rdd.RDD
object Demo02_flatMap_wordCount {
def main(args: Array[String]): Unit = {
val sc: SparkContext = ScUtils.sparkContext();
val rdd: RDD[String] = sc.textFile("file:///D:/a/a.txt");
//可以直接使用countByKey实现统计
rdd.flatMap(_.split("\\s+")).map((_, 1)).countByKey().foreach(kv => {
println(kv._1 + "\t" + kv._2);
});
sc.stop();
}
}
以下是输出的结果:
Alex1
Rose1
Hello4
Jack1
Mary1
3:mapPartitions算子
def mapPartitions[U : ClassTag](f: scala.Iterator[T] => scala.Iterator[U],
preservesPartitioning: Boolean = false): RDD[U]
该函数和map函数类似,只不过映射函数的参数由RDD中的每一个元素变成了RDD中每一个分区的迭代器。如果在映射的过程中需要频繁创建额外的对象,使用mapPartitions要比map高效的多。
比如,将RDD中的所有数据通过JDBC连接写入数据库,如果使用map函数,可能要为每一个元素都创建一个connection,这样开销很大,如果使用mapPartitions,那么只需要针对每一个分区建立一个connection。
以下代码,定义一个数组,分成两个区,计算每一个区数字的和。
package cn.wang
import cn.wang.utils.ScUtils
import org.apache.spark.SparkContext
import org.apache.spark.rdd.RDD
object Demo02_mapPartitions {
def main(args: Array[String]): Unit = {
val sc: SparkContext = ScUtils.sparkContext();
//声明一个数组,分两个区
val rdd: RDD[Int] = sc.parallelize(Array(1, 2, 3, 4, 5, 6), 2);
//以下按分区对数据进行合并求和
val rdd2: RDD[Int] = rdd.mapPartitions(p => {
println("分区被调用");
var list = List[Int]();
var i: Int = 0;
while (p.hasNext) {
i += p.next();
}
println("每一个分区计算完成以后,i的值为:" + i);
//mapPartitions返回的是一个Iterator所以,需要将List转成Iterator对象
(list.::(i)).toIterator;
});
rdd2.collect().foreach(println);
sc.stop();
}
}
输出的结果为:
6
15
上面代码计算的图示为:
第一个分区: 1,2 ,3 。相加,结果为6.
第二个分区:4,5,6。相加结果为15。返回两个数组
以下在命令行中,执行,结果相同。
scala> val rdd = sc.parallelize(1 to 6,2);
scala> val rdd2 = rdd.mapPartitions(arr=>{
| var list = List[Int]();
| var i:Int=0;
| while(arr.hasNext){
| i+=arr.next();
| }
| (list.::(i)).toIterator;
| });
rdd2: org.apache.spark.rdd.RDD[Int] = MapPartitionsRDD[2] at mapPartitions at <console>:26
scala> rdd2.collect();
res2: Array[Int] = Array(6, 15)
4:mapPartitionsWithIndex算子
源代码:
def mapPartitionsWithIndex[U : ClassTag](f: (Int, scala.Iterator[T]) => scala.Iterator[U],
preservesPartitioning: Boolean = false): RDD[U]
参数:
1:是分区的下标。从0开始。
2:被迭代的对象。
返回:
一个新的RDD[Iterator]
以下示例,将返回一个按分区显示的Map集合:
要求有源数据:(1,2,3,4,5,6),2 : 分两个区,现在将它们按区的ID将数据组织成:
(0,List(1,2,3));
(1,List(4,5,6));
示例代码如下:
package cn.wang
import cn.wang.utils.ScUtils
import org.apache.spark.SparkContext
import org.apache.spark.rdd.RDD
import scala.collection.mutable
object Demo03_mapPartitionsWithIndex {
def main(args: Array[String]): Unit = {
val sc: SparkContext = ScUtils.sparkContext();
//声明一个数组,分两个区
val rdd: RDD[Int] = sc.parallelize(1 to 6, 2);
//使用mapPartitionsWithIndex进行分区计算
val rdd2 = rdd.mapPartitionsWithIndex((index: Int, arr) => {
println("下标:" + index);//查看下标值
//使用HashMap声明一个可变的集合,如果直接使用 Map则为不可变的集合
var map = new mutable.HashMap[Int, List[Int]]();
//声明一个List用于追加数据
var list = List[Int]();
//对每一个区进行迭代
while (arr.hasNext) {
//获取某一个区中,的某一个数据
val num = arr.next();
//输出某个分区下的某个数值
println("下标:" + index + ",数值:" + num);
//将 个数添加到List集合中,注意要再赋给List对象
//因为RDD为不可变的数据集,所以,在追加数据以后必须要再赋给这个RDD
//否则里面会没有数据
list = list.:+(num);//注意赋值
}
//声明的map是一个可变的集合,所以可以向里面追加数据
map.put(index, list);
//返回这个集合对象
map.iterator;
});
//输出数据
rdd2.collect().foreach(println(_));
sc.stop();
}
}
输出结果:
(0,List(1, 2, 3))
(1,List(4, 5, 6))
你也可以修改数据输出的格式,修改rdd2.collect().foreach(..)为以下代码:
rdd2.collect().foreach(entry => {
print("分区下标:" + entry._1);
print("\t分区数据为:");
println(entry._2.mkString(",")); //将数据组使用,组成字符串
});
则输出的结果为:
分区下标:0分区数据为:1,2,3
分区下标:1分区数据为:4,5,6
5:glom算子
源代码:
def glom(): RDD[Array[T]] = withScope {
new MapPartitionsRDD[Array[T], T](this, (context, pid, iter) => Iterator(iter.toArray))
}
功能:
将不同分区中的数据T,转换放到不同的Array[T]中。
示例1:
scala> var rdd1 = sc.makeRDD(1 to 10,3);
scala> rdd1.collect;
res32: Array[Int] = Array(1, 2, 3, 4, 5, 6, 7, 8, 9, 10)
scala> val rdd2 = rdd1.glom();
rdd2: org.apache.spark.rdd.RDD[Array[Int]] = MapPartitionsRDD[22] at glom at <console>:26
scala> rdd2.collect();
res33: Array[Array[Int]] = Array(Array(1, 2, 3), Array(4, 5, 6), Array(7, 8, 9, 10))
示例2:
scala> val rdd = sc.makeRDD(List(("A",1),("B",2),("C",3)),2);
rdd: org.apache.spark.rdd.RDD[(String, Int)] = ParallelCollectionRDD[23] at makeRDD at <console>:24
scala> rdd.collect
res35: Array[(String, Int)] = Array((A,1), (B,2), (C,3))
scala> val rdd2 = rdd.glom();
rdd2: org.apache.spark.rdd.RDD[Array[(String, Int)]] = MapPartitionsRDD[24] at glom at <console>:26
scala> rdd2.collect
res34: Array[Array[(String, Int)]] = Array(Array((A,1)), Array((B,2), (C,3)))
6:randomSplit算子
源码:
def randomSplit(weights: Array[Double],
seed: Long = Utils.random.nextLon...): Array[RDD[T]]
参数:
Array数组。用于指定权重,权重越高,越可能多的获取更多的元素。
返回:
返回RDD的数组即:Array[RDD[T]]
示例:
scala> val rdd = sc.makeRDD(1 to 10 ,2);
rdd: org.apache.spark.rdd.RDD[Int] = ParallelCollectionRDD[25] at makeRDD at <console>:24
//以下指定返回两个RDD,第二个RDD的权重为3.0,第一个为1.0,即第二个RDD会获取更多的元素个数
scala> val array = rdd.randomSplit(Array(1.0,3.0));
array: Array[org.apache.spark.rdd.RDD[Int]] = Array(MapPartitionsRDD[26] at randomSplit at <console>:26, MapPartitionsRDD[27] at randomSplit at <console>:26)
scala> array.length;
res36: Int = 2
scala> array(0).collect;
res37: Array[Int] = Array(2, 5, 10)
scala> array(1).collect;
res38: Array[Int] = Array(1, 3, 4, 6, 7, 8, 9)
7、zip算子
用于将两个算子按key-value对组成新的算子。如果两个RDD的分区不一样,则会直接抛出异常。且每一个分区中的元素的个数,也必须要一样,否则一样也会抛出异常。
源代码:
/**
* Zips this RDD with another one, returning key-value pairs with the first element in each RDD,
* second element in each RDD, etc. Assumes that the two RDDs have the *same number of
* partitions* and the *same number of elements in each partition* (e.g. one was made through
* a map on the other).
*/
def zip[U: ClassTag](other: RDD[U]): RDD[(T, U)] = withScope {
zipPartitions(other, preservesPartitioning = false) { (thisIter, otherIter) =>
new Iterator[(T, U)] {
def hasNext: Boolean = (thisIter.hasNext, otherIter.hasNext) match {
case (true, true) => true
case (false, false) => false
case _ => throw new SparkException("Can only zip RDDs with " +
"same number of elements in each partition")
}
def next(): (T, U) = (thisIter.next(), otherIter.next())
}
}
}
示例1
package cn.wang
import cn.wang.utils.ScUtils
import org.apache.spark.SparkContext
import org.apache.spark.rdd.RDD
object Demo02_zip {
def main(args: Array[String]): Unit = {
val sc: SparkContext = ScUtils.sparkContext();
val rdd1: RDD[Int] = sc.parallelize(Seq(1, 2, 3, 4), 1); // 声明的集合,声明几个分区就是几个分区,只有文件才会去计算分区数量
val rdd2: RDD[String] = sc.parallelize(Seq("A", "B", "C", "D"), 1);
val rdd3: RDD[(Int, String)] = rdd1.zip(rdd2);
println("*******合并以后的结果**********");
rdd3.collect().foreach(kv => {
println(kv._1 + "," + kv._2);
});
sc.stop();
}
}
结果:
1,A
2,B
3,C
4,D
图示:

8、zipParititions算子

org.apache.spark.rdd.RDD
def zipPartitions[B : ClassTag, V : ClassTag](rdd2: RDD[B],
preservesPartitioning: Boolean)
(f: (scala.Iterator[T], scala.Iterator[B]) => scala.Iterator[V]): RDD[V]
Zip this RDD's partitions with one (or more) RDD(s) and return a new RDD by applying a function to the zipped partitions. Assumes that all the RDDs have the *same number of partitions*, but does *not* require them to have the same number of elements in each partition.
用于将一个或是多个RDD进行合并,被合并的RDD必须要拥有相同的分区数量,但每一个分区中的元素,可以不同。
示例:
package cn.wang
import cn.wang.utils.ScUtils
import org.apache.spark.SparkContext
import org.apache.spark.rdd.RDD
object Demo02_zipPartitions {
def main(args: Array[String]): Unit = {
val sc: SparkContext = ScUtils.sparkContext();
val rdd1: RDD[Int] = sc.makeRDD(1 to 10, 4);
val rdd2: RDD[String] = sc.makeRDD(Seq("A", "B", "C", "D"), 4);
//zipPartitions必须返回一个Iterator集合
val rdd3: RDD[String] = rdd1.zipPartitions(rdd2)((iter1, iter2) => {
var list = List[String]();
while (iter1.hasNext && iter2.hasNext) {
//::即在最前面追加元素
list = list :+ (iter1.next() + "_" + iter2.next());
}
list.iterator;
});
println(rdd3.count());
println("*******输出的数据是***********")
rdd3.foreach(str => {
println(str);
});
sc.stop();
}
}
结果:
3_B
1_A
8_D
6_C
图示:





