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

Spark-4 RDD算子-1

Coding On Road 2017-12-24
169

【郑重声明:本文为原创作品,版权所有,如转载请注明出处。

         如作为商业宣传,本人将保留诉讼的权力。--王健】

  同时欢迎技术爱好者提出问题,共同进步。本文示展示的程序,都是用Scala语言编写。


【郑重声明:本文为原创作品,版权所有,如转载请注明出处。

         如作为商业宣传,本人将保留诉讼的权力。】


Spark RDD编程

 

1RDD的概念

RDD(Resilient Distributed Datasets) ,弹性分布式数据集, 是分布式内存的一个抽象概念,RDD提供了一种高度受限的共享内存模型,即RDD是只读的记录分区的集合,只能通过在其他RDD执行确定的转换操作(如mapjoingroup 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是个抽象类,具体由各子类实现,如MappedRDDShuffledRDD等子类。 Spark将常用的大数据操作都转化成为RDD的子类。

 

4、RDD算子

Spark算子是由Scala写成的,所以,看懂Scala语法很重要。在讲解以下算子时,本人会附上大量的Scala的源代码,帮助大家理解。

 

1、转换算子

转换算子,生成新的RDD。不触发Job。但有些转换算子如repartition会有shuffle操作,但依然不会执行job。只有行动算子,才会触发Job

1、输入分区与输出分区一对一类型

即输出时的分区与原分区保持相同:

//设置为分7个分区如果文件正好7个字节的话,否则将会根据算法,重新计算分区//数量,见HadoopFileInputFormatgetSplits方法源代码(旧版),因到到作者写此文//章时,spark2.1.2依然使用的是Hadoopmapred下的类,即旧版本的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算子最后最后的返回值来决定。

 

 

2flatMap算子

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:单词统计

 

以下示例,用flatMapmap的组合实现单词的统计:

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

 

 

 

3mapPartitions算子

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

上面代码计算的图示为:

 

第一个分区: 12 3 。相加,结果为6.

第二个分区:456。相加结果为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)

 

 

 

4mapPartitionsWithIndex算子

源代码:

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

 

5glom算子

源代码:

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)))

 

 

6randomSplit算子

源码:

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

 

图示:


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

评论