RDD持久化
1. RDD Cache 快取
說明
RDD 通過Cache 或者Persist 方法將前面的計算結果快取,默認情況下會把資料以快取在JVM 的堆記憶體中,但是并不是這兩個方法被呼叫時立即快取,而是觸發后面的 action 算子時,該RDD 將會被快取在計算節點的記憶體中,并供后面重用,
// cache 操作會增加血緣關系,不改變原有的血緣關系
println(wordToOneRdd.toDebugString)
// 資料快取,
wordToOneRdd.cache()
// 可以更改存盤級別
//mapRdd.persist(StorageLevel.MEMORY_AND_DISK_2)
案例實操
package com.atguigu.bigdata.spark.core.rdd.Persist
import org.apache.spark.rdd.RDD
import org.apache.spark.{SparkConf, SparkContext}
object Spark03_RDD_Persist {
def main(args: Array[String]): Unit = {
val sparkConf: SparkConf = new SparkConf().setMaster("local[*]").setAppName("Operator")
val sc = new SparkContext(sparkConf)
val list = List("Hello Spark","Hello Scala")
val rdd = sc.makeRDD(list)
val flatRDD = rdd.flatMap(_.split(" "))
val mapRDD = flatRDD.map(
word => {
println("@@@@@@@@@@")
(word,1)
}
)
mapRDD.cache() ///持久化操作
val reduceRDD: RDD[(String, Int)] = mapRDD.reduceByKey(_ + _)
reduceRDD.collect().foreach(println)
println("***************************")
val groupRDD: RDD[(String, Iterable[Int])] = mapRDD.groupByKey()
groupRDD.collect().foreach(println)
sc.stop()
}
}
存盤級別
object StorageLevel {
val NONE = new StorageLevel(false, false, false, false)
val DISK_ONLY = new StorageLevel(true, false, false, false)
val DISK_ONLY_2 = new StorageLevel(true, false, false, false, 2)
val MEMORY_ONLY = new StorageLevel(false, true, false, true)
val MEMORY_ONLY_2 = new StorageLevel(false, true, false, true, 2)
val MEMORY_ONLY_SER = new StorageLevel(false, true, false, false)
val MEMORY_ONLY_SER_2 = new StorageLevel(false, true, false, false, 2)
val MEMORY_AND_DISK = new StorageLevel(true, true, false, true)
val MEMORY_AND_DISK_2 = new StorageLevel(true, true, false, true, 2)
val MEMORY_AND_DISK_SER = new StorageLevel(true, true, false, false)
val MEMORY_AND_DISK_SER_2 = new StorageLevel(true, true, false, false, 2)
val OFF_HEAP = new StorageLevel(true, true, true, false, 1)

快取有可能丟失,或者存盤于記憶體的資料由于記憶體不足而被洗掉,RDD 的快取容錯機制保證了即使快取丟失也能保證計算的正確執行,通過基于RDD 的一系列轉換,丟失的資料會被重算,由于RDD 的各個Partition 是相對獨立的,因此只需要計算丟失的部分即可,并不需要重算全部Partition,
Spark 會自動對一些 Shuffle 操作的中間資料做持久化操作(比如:reduceByKey),這樣做的目的是為了當一個節點 Shuffle 失敗了避免重新計算整個輸入,但是,在實際使用的時候,如果想重用資料,仍然建議呼叫persist 或 cache,
2. RDD CheckPoint 檢查點
說明
所謂的檢查點其實就是通過將RDD 中間結果寫入磁盤,
由于血緣依賴過長會造成容錯成本過高,這樣就不如在中間階段做檢查點容錯,如果檢查點之后有節點出現問題,可以從檢查點開始重做血緣,減少了開銷,
對RDD 進行checkpoint 操作并不會馬上被執行,必須執行Action 操作才能觸發,
案例實操
package com.atguigu.bigdata.spark.core.rdd.Persist
import org.apache.spark.rdd.RDD
import org.apache.spark.{SparkConf, SparkContext}
object Spark04_RDD_Persist {
def main(args: Array[String]): Unit = {
val sparkConf: SparkConf = new SparkConf().setMaster("local[*]").setAppName("Operator")
val sc = new SparkContext(sparkConf)
sc.setCheckpointDir("cp")
val list = List("Hello Spark","Hello Scala")
val rdd = sc.makeRDD(list)
val flatRDD = rdd.flatMap(_.split(" "))
val mapRDD = flatRDD.map(
word => {
println("@@@@@@@@@@")
(word,1)
}
)
//checkpoint需要落盤,需要指定檢查點保存的路徑
//檢查點保存的檔案,作業執行完,不會洗掉
//一般的保存路徑都是在分布式存盤系統中,HDFS
mapRDD.checkpoint()
val reduceRDD: RDD[(String, Int)] = mapRDD.reduceByKey(_ + _)
reduceRDD.collect().foreach(println)
println("***************************")
val groupRDD: RDD[(String, Iterable[Int])] = mapRDD.groupByKey()
groupRDD.collect().foreach(println)
sc.stop()
}
}
3. 快取和檢查點區別
(1)Cache 快取只是將資料臨時保存起來進行重用,不切斷血緣依賴,它會在血緣關系中添加新的依賴,一旦出現問題,它可以從頭讀取資料,Checkpoint 檢查點切斷血緣依賴,會重新建立新的血緣關系,它等同于改變的資料源,同時將資料長久的保存在磁盤檔案中進行資料重用,
(2)Cache 快取的資料通常存盤在磁盤、記憶體等地方,可靠性低,Checkpoint 的資料通常存盤在HDFS 等容錯、高可用的檔案系統,可靠性高,
(3)建議對checkpoint () 的RDD 使用Cache 快取,這樣 checkpoint 的job 只需從 Cache 快取
中讀取資料即可,否則需要再從頭計算一次RDD,
(4)persist:將資料臨時存盤在磁盤檔案中進行資料重用,涉及到磁盤IO,性能較低,但是資料安全;如果作業執行完畢,臨時保存的資料檔案就會丟失,
轉載請註明出處,本文鏈接:https://www.uj5u.com/qita/295361.html
標籤:其他
上一篇:mycat入門
