4.1
RDD持久化原理
4.2
持久化策略
4.3
如何选择RDD持久化策略
4.4
共享变量的工作原理

4.1
RDD持久化原理
Spark非常重要的一个功能特性就是可以将RDD持久化在内存中。当对RDD执行持久化操作时,每个节点都会将自己操作的RDD的partition持久化到内存中,并且在之后对该RDD的反复使用中,直接使用内存缓存的partition。这样的话,对于针对一个RDD反复执行多个操作的场景,就只要对RDD计算一次即可,后面可直接使用该RDD,而不需要反复计算多次该RDD。
巧妙使用RDD持久化,甚至在某些场景下,可以将spark应用程序的性能提升10倍。对于迭代式算法和快速交互式应用来说,RDD持久化,是非常重要的。
要持久化一个RDD,只要调用cache()或者persist()方法即可。在该RDD第一次被计算出来时,就会直接缓存在每个节点中。而且Spark的持久化机制还是自动容错的,如果持久化的RDD的任何partition丢失了,那么spark会自动通过其源RDD,使用transformation操作重新计算该partition。
cache()和persist()的区别在于,cache()是persist()的一种简化方式,cache()的底层就是调用的persist()的无参版本,同时就是调用persist(MEMORY_ONLY),将数据持久化到内存中。如果需要从内存中清除缓存,那么可以使用unpersist()方法。
不使用持久化:

使用持久化:


4.2
持久化策略
持久化策略 | |
MEMORY_ONLY | 以非序列化的Java对象的方式持久化在JVM内存中,如果内存无法完全存储RDD所有的partition,那么那些没有持久化的partition就会在下一次需要使用它的时候,重新被计算 |
MEMORY_AND_DISK | 同上,但是当某些partition无法存储在内存中时,会持久化到磁盘中。下次需要使用这些partition时,需要从磁盘上读取 |
MEMORY_ONLY_SER | 同MEMORY_ONLY,但是会使用Java序列化方式,将Java对象序列化后进行持久化。可以减少内存开销,但是需要进行反序列化,因此会加大CPU开销 |
MEMORY_AND_DSK_SER | 同MEMORY_AND_DSK。但是使用序列化方式持久化Java对象 |
DISK_ONLY | 使用非序列化Java对象的方式持久化,完全存储到磁盘上 |
MEMORY_ONLY_2 | —— |
MEMORY_AND_DISK_2 | 如果是尾部加了2的持久化级别,表示会将持久化数据复用一份,保存到其他节点,从而在数据丢失时,不需要再次计算,只需要使用备份数据即可 |

4.3
如何选择RDD持久化策略
Spark提供的多种持久化级别,主要是为了在CPU和内存消耗之间进行取舍。
下面是一些通用的持久化级别的选择建议:
优先使用MEMORY_ONLY,如果可以缓存所有数据的话,那么就使用这种策略。因为纯内存速度最快,而且没有序列化,不需要消耗CPU进行反序列化操作。
如果MEMORY_ONLY策略,无法存储的下所有数据的话,那么使用MEMORY_ONLY_SER,将数据进行序列化进行存储,纯内存操作还是非常快,只是要消耗CPU进行反序列化。
如果需要进行快速的失败恢复,那么就选择带后缀为_2的策略,进行数据的备份,这样在失败时,就不需要重新计算了。
能不使用DISK相关的策略,就不用使用,有的时候,从磁盘读取数据,还不如重新计算一次。
案例实战:计算不使用持久化和使用持久化读取文件的时间差
import org.apache.spark.SparkConf;import org.apache.spark.api.java.JavaRDD;import org.apache.spark.api.java.JavaSparkContext;public class persistJava {public static JavaSparkContext getsc(){SparkConf sparkConf = new SparkConf().setAppName("parallelize").setMaster("local");JavaSparkContext sc = new JavaSparkContext(sparkConf);return sc;}public static void main(String[] args) {JavaRDD lines = getsc().textFile("hdfs://bigdata-01.com:9000/user/zjy/datas/test");//JavaRDD lines = getsc().textFile("hdfs://bigdata-01.com:9000/user/zjy/datas/test").cache(); 使用持久化long begin = System.currentTimeMillis();System.out.println(lines.count());System.out.println(System.currentTimeMillis() - begin); //第一次使用的时间long begin1 = System.currentTimeMillis();System.out.println(lines.count());System.out.println(System.currentTimeMillis() - begin1); //第二次使用的时间}}
打印结果不使用持久化的时间:

打印结果使用持久化的时间:

可以看到时间还是有所区别的,当然本文采用的分布式文件size并不是很大。如果是几十m的文件,时间相差可能会有十几倍。

4.4
共享变量的工作原理
Spark一个非常重要的特性就是共享变量。
默认情况下,如果在一个算子的函数中使用到了某个外部的变量,那么这个变量的值会被拷贝到每个task中。此时每个task只能操作自己的那份变量副本。如果多个task想要共享某个变量,那么这种方式是做不到的。
Spark为此提供了两种共享变量,一种是Broadcast Variable(广播变量),另一种是Accumulator (累加变量)。Broadeast Variable会将使用到的变量,仅仅为每个节点拷贝一份,更大的用处是优化性能,减少网络传输以及内存消耗。Accumulator则可以让多个task共同操作一份变量,主要可以进行累加操作。

默认情况下,算子的函数内使用到的外部变量时,会拷贝到执行这个函数的每一个task中。
当变量比较大的时候,网络传输的量就会大,并且在每个节点上占用较多的内存空间。
不使用共享变量:

如果把算子使用的变量设置成共享变量的话,那么变量只会拷贝一份到每一个worker节点上,节点上所有的task都会共享这个一份变量。
使用共享变量:


Broadcast Variable(广播变量):
Spark提供的Broadcast Variable,是只读的。并且在每个节点上只会有一份副本,而不会为每个task都拷贝一份副本。因此其最大作用,就是减少变量到各个节点的网络传输消耗,以及在各个节点上的内存消耗。此外,spark自己内部也使用了高效的广播算法来减少网络消耗。
可以通过调用SparkContext的broadcast()方法,来针对某个变量创建广播变量。然后在算子的函数内,使用到广播变量时,每个节点只会拷贝一份副本了。每个节点可以使用广播变量的value0方法获取值。
记住,广播变量,是只读的。
案例实战:通过两种不同的广播方式,将list的值进行乘积加和。
//JAVAimport org.apache.spark.Accumulator;import org.apache.spark.SparkConf;import org.apache.spark.api.java.JavaRDD;import org.apache.spark.api.java.JavaSparkContext;import org.apache.spark.api.java.function.Function;import org.apache.spark.api.java.function.VoidFunction;import org.apache.spark.broadcast.Broadcast;import java.util.Arrays;import java.util.List;public class VarJava {public static JavaSparkContext getsc() {SparkConf sparkconf = new SparkConf().setAppName("parallelize").setMaster("local");JavaSparkContext sc = new JavaSparkContext(sparkconf);return sc;}public static void main(String[] args) {List list = Arrays.asList(1,2,3,4);JavaSparkContext sc = getsc();final Broadcast broadcast = sc.broadcast(10); //Broadcast Variable(广播变量)JavaRDD rdd = sc.parallelize(list);JavaRDD values = rdd.map(new Function<Integer,Integer>() {@Overridepublic Integer call(Integer value) throws Exception {return value * (int)broadcast.getValue(); //list的值*10}});values.foreach(new VoidFunction() {@Overridepublic void call(Object o) throws Exception {System.out.println(o);}});final Accumulator<Integer> accumulator = sc.accumulator(5); //Accumulator (累加变量)rdd.foreach(new VoidFunction<Integer>() {@Overridepublic void call(Integer o) throws Exception {accumulator.add(o); //将list的值加和}});System.out.println(accumulator.value());}}
//Scalaimport org.apache.spark.{SparkConf, SparkContext}object VarScala {def getsc(): SparkContext = {val sparkconf = new SparkConf().setAppName("parallelize").setMaster("local")val sc = new SparkContext(sparkconf)return sc}def main(args: Array[String]): Unit = {val list = Array(1,2,3,4)val sc =getsc()val rdd = sc.parallelize(list)val broadcast = sc.broadcast(10) //Broadcast Variable(广播变量)val values = rdd.map(x => x*broadcast.value)values.foreach(x => System.out.println(x))val accumulator = sc.longAccumulator //Accumulator (累加变量)val value = rdd.foreach(x => accumulator.add(x))System.out.println(accumulator.value)}}
打印结果:


接下来就要应用前几章所讲到的算子进行一些简单和复杂的排序操作。
欢迎收看下一章,大数据分析之Spark核心编程(五)。






