《大數(shù)據(jù)離線分析技術(shù)》課件-項(xiàng)目5:RDD核心概念與操作_第1頁(yè)
《大數(shù)據(jù)離線分析技術(shù)》課件-項(xiàng)目5:RDD核心概念與操作_第2頁(yè)
《大數(shù)據(jù)離線分析技術(shù)》課件-項(xiàng)目5:RDD核心概念與操作_第3頁(yè)
《大數(shù)據(jù)離線分析技術(shù)》課件-項(xiàng)目5:RDD核心概念與操作_第4頁(yè)
《大數(shù)據(jù)離線分析技術(shù)》課件-項(xiàng)目5:RDD核心概念與操作_第5頁(yè)
已閱讀5頁(yè),還剩88頁(yè)未讀 繼續(xù)免費(fèi)閱讀

付費(fèi)下載

下載本文檔

版權(quán)說(shuō)明:本文檔由用戶提供并上傳,收益歸屬內(nèi)容提供方,若內(nèi)容存在侵權(quán),請(qǐng)進(jìn)行舉報(bào)或認(rèn)領(lǐng)

文檔簡(jiǎn)介

RDD概述什么是RDDRDD的屬性什么是RDD/01

RDD(ResilientDistributedDataset)叫做彈性分布式數(shù)據(jù)集,是Spark中最基本的數(shù)據(jù)抽象。代碼中是一個(gè)抽象類,它代表一個(gè)彈性的、不可變、可分區(qū)、里面的元素可并行計(jì)算的集合。什么是RDDRDD中的彈性是指:存儲(chǔ)的彈性:內(nèi)存與磁盤的自動(dòng)切換。容錯(cuò)的彈性:數(shù)據(jù)丟失可以自動(dòng)恢復(fù)。計(jì)算的彈性:計(jì)算出錯(cuò)重試機(jī)制。分片的彈性:可根據(jù)需要重新分片。什么是RDD

RDD中的不可變是指創(chuàng)建一個(gè)RDD如果更改,并不是真正意義上的更改,只是又創(chuàng)建了一個(gè)新的RDD。RDD中可分區(qū)是指能進(jìn)行分區(qū)。RDD中并行計(jì)算是指,因?yàn)镽DD的分區(qū)特性,所以支持并行處理的特性。即不同節(jié)點(diǎn)上的數(shù)據(jù)可以分別被處理,然后生成一個(gè)新的RDD。什么是RDD

RDD的結(jié)構(gòu)如圖所示:什么是RDDRDD的屬性/02RDD的屬性如下:一組分區(qū)(Partition),即數(shù)據(jù)集的基本組成單位;一個(gè)計(jì)算每個(gè)分區(qū)的函數(shù);RDD之間的依賴關(guān)系;一個(gè)Partitioner,即RDD的分片函數(shù);一個(gè)列表,存儲(chǔ)存取每個(gè)Partition的優(yōu)先位置(preferredlocation)。RDD的屬性1.什么是RDD

2.RDD的屬性

RDD的特點(diǎn)RDD特點(diǎn)概述RDD的特點(diǎn)詳解RDD特點(diǎn)概述/01

RDD表示只讀的分區(qū)的數(shù)據(jù)集,對(duì)RDD進(jìn)行改動(dòng),只能通過(guò)RDD的轉(zhuǎn)換操作,由一個(gè)RDD得到一個(gè)新的RDD,新的RDD包含了從其他RDD衍生所必需的信息。RDDs之間存在依賴,RDD的執(zhí)行是按照血緣關(guān)系延時(shí)計(jì)算的。如果血緣關(guān)系較長(zhǎng),可以通過(guò)持久化RDD來(lái)切斷血緣關(guān)系。RDD特點(diǎn)概述RDD的特點(diǎn)詳解/02

存儲(chǔ)的彈性:內(nèi)存與磁盤的自動(dòng)切換;容錯(cuò)的彈性:數(shù)據(jù)丟失可以自動(dòng)恢復(fù);計(jì)算的彈性:計(jì)算出錯(cuò)重試機(jī)制;分片的彈性:可根據(jù)需要重新分片。RDD的特點(diǎn)詳解-彈性RDD邏輯上是分區(qū)的,每個(gè)分區(qū)的數(shù)據(jù)是抽象存在的,計(jì)算時(shí)通過(guò)compute函數(shù)得到每個(gè)分區(qū)的數(shù)據(jù)。如果RDD是通過(guò)已有的文件系統(tǒng)構(gòu)建,則compute函數(shù)是讀取指定文件系統(tǒng)中的數(shù)據(jù),如果RDD是通過(guò)其他RDD轉(zhuǎn)換而來(lái),則compute函數(shù)是執(zhí)行轉(zhuǎn)換邏輯將其他RDD的數(shù)據(jù)進(jìn)行轉(zhuǎn)換。RDD的特點(diǎn)詳解-分區(qū)由一個(gè)RDD轉(zhuǎn)換到另一個(gè)RDD,可以通過(guò)算子實(shí)現(xiàn)。RDD的操作算子包括兩類,一類叫做transformations,它是用來(lái)將RDD進(jìn)行轉(zhuǎn)化,構(gòu)建RDD的血緣關(guān)系;另一類叫做actions,它是用來(lái)觸發(fā)RDD的計(jì)算,得到RDD的相關(guān)計(jì)算結(jié)果或者將RDD保存的文件系統(tǒng)中。RDD的特點(diǎn)詳解-只讀RDDs通過(guò)操作算子進(jìn)行轉(zhuǎn)換,轉(zhuǎn)換得到的新RDD包含了從其他RDDs衍生所必需的信息,RDDs之間維護(hù)著這種血緣關(guān)系,也稱之為依賴。如下圖所示,依賴包括兩種,一種是窄依賴,RDDs之間分區(qū)是一一對(duì)應(yīng)的,另一種是寬依賴,下游RDD的每個(gè)分區(qū)與上游RDD的每個(gè)分區(qū)都有關(guān),是多對(duì)多的關(guān)系。RDD的特點(diǎn)詳解-依賴如果在應(yīng)用程序中多次使用同一個(gè)RDD,可以將該RDD緩存起來(lái),該RDD只有在第一次計(jì)算的時(shí)候會(huì)根據(jù)血緣關(guān)系得到分區(qū)的數(shù)據(jù),在后續(xù)其他地方用到該RDD的時(shí)候,會(huì)直接從緩存處取而不用再根據(jù)血緣關(guān)系計(jì)算,這樣就加速后期的重用。

RDD的特點(diǎn)詳解-緩存如下圖所示,RDD-1經(jīng)過(guò)一系列的轉(zhuǎn)換后得到RDD-n并保存到hdfs,RDD-1在這一過(guò)程中會(huì)有個(gè)中間結(jié)果,如果將其緩存到內(nèi)存,那么在隨后的RDD-1轉(zhuǎn)換到RDD-m這一過(guò)程中,就不會(huì)計(jì)算其之前的RDD-0。RDD的特點(diǎn)詳解-緩存1.RDD特點(diǎn)概述

2.RDD的特點(diǎn)詳解

RDD的創(chuàng)建創(chuàng)建RDD的三種方式創(chuàng)建方式詳解創(chuàng)建RDD的三種方式/01

在RDD中,通常就代表和包含了Spark應(yīng)用程序的輸入源數(shù)據(jù)。在創(chuàng)建了初始的RDD之后,才可以通過(guò)SparkCore提供的transformation算子,對(duì)該RDD進(jìn)行transformation(轉(zhuǎn)換)操作,來(lái)獲取其他的RDD。創(chuàng)建RDD的三種方式SparkCore為我們提供了三種創(chuàng)建RDD的方式,包括:1.使用程序中的集合創(chuàng)建RDD2.使用本地文件創(chuàng)建RDD3.使用HDFS文件創(chuàng)建RDD創(chuàng)建RDD的三種方式創(chuàng)建方式詳解/02

1、使用程序中的集合創(chuàng)建RDD,主要用于進(jìn)行測(cè)試,可以在實(shí)際部署到集群運(yùn)行之前,使用集合構(gòu)造測(cè)試數(shù)據(jù),來(lái)測(cè)試后面的spark應(yīng)用的流程。

創(chuàng)建方式詳解2、使用本地文件創(chuàng)建RDD,主要用于臨時(shí)性地處理一些存儲(chǔ)了大量數(shù)據(jù)的文件。集群上運(yùn)行時(shí)需要所有集群上都有該文件。

創(chuàng)建方式詳解3、使用HDFS文件創(chuàng)建RDD,是最常用的生產(chǎn)環(huán)境處理方式,主要可以針對(duì)HDFS上存儲(chǔ)的大數(shù)據(jù),進(jìn)行離線批處理操作。

創(chuàng)建方式詳解1.創(chuàng)建RDD的三種方式

2.創(chuàng)建方式詳解

RDD的分區(qū)RDD分區(qū)介紹RDD分區(qū)原則RDD設(shè)置分區(qū)的方法RDD分區(qū)介紹/01

RDD是彈性分布式數(shù)據(jù)集,通常RDD很大,會(huì)被分成很多個(gè)分區(qū),分別保存在不同的節(jié)點(diǎn)上,作用有二:增加并行度和減少通信開銷(連接操作),例如下圖:

RDD分區(qū)介紹RDD分區(qū)原則/02

RDD分區(qū)的一個(gè)原則是使得分區(qū)的個(gè)數(shù)盡量等于集群中的CPU核心(core)數(shù)目。對(duì)于不同的Spark部署模式而言(本地模式、Standalone模式、YARN模式、Mesos模式),都可以通過(guò)設(shè)置spark.default.parallelism這個(gè)參數(shù)的值,來(lái)配置默認(rèn)的分區(qū)數(shù)目。RDD分區(qū)原則*本地模式:默認(rèn)為本地機(jī)器的CPU數(shù)目,若設(shè)置了local[N],則默認(rèn)為N,local[*]則自動(dòng)判斷。*ApacheMesos:默認(rèn)的分區(qū)數(shù)為8。*Standalone或YARN:在“集群中所有CPU核心數(shù)目總和”和“2”二者中取較大值作為默認(rèn)值。RDD分區(qū)原則RDD設(shè)置分區(qū)方法/03

創(chuàng)建RDD時(shí)手動(dòng)指定分區(qū)個(gè)數(shù)。使用reparititon方法重新設(shè)置分區(qū)個(gè)數(shù)。自定義分區(qū)方法:Spark提供了自帶的HashPartitioner(哈希分區(qū))與RangePartitioner(區(qū)域分區(qū)),能夠滿足大多數(shù)應(yīng)用場(chǎng)景的需求。與此同時(shí),Spark也支持自定義分區(qū)方式,即通過(guò)提供一個(gè)自定義的Partitioner對(duì)象來(lái)控制RDD的分區(qū)方式,從而利用領(lǐng)域知識(shí)進(jìn)一步減少通信開銷。RDD設(shè)置分區(qū)方法1.RDD分區(qū)介紹

2.RDD分區(qū)原則

3.RDD設(shè)置分區(qū)方法

窄依賴窄依賴介紹窄依賴對(duì)優(yōu)化的優(yōu)勢(shì)窄依賴介紹/01

窄依賴是指1個(gè)父RDD分區(qū)對(duì)應(yīng)1個(gè)子RDD的分區(qū)。換句話說(shuō),一個(gè)父RDD的分區(qū)對(duì)應(yīng)于一個(gè)子RDD的分區(qū),或者多個(gè)父RDD的分區(qū)對(duì)應(yīng)于一個(gè)子RDD的分區(qū)。窄依賴介紹

窄依賴分為兩種情況:1個(gè)子RDD的分區(qū)對(duì)應(yīng)于1個(gè)父RDD的分區(qū),比如map,filter,union等算子。1個(gè)子RDD的分區(qū)對(duì)應(yīng)于N個(gè)父RDD的分區(qū),比如co-partionedjoin。窄依賴介紹

窄依賴可以支持在同一個(gè)集群Executor上,以pipeline管道形式順序執(zhí)行多條命令,例如在執(zhí)行了map后,緊接著執(zhí)行filter。分區(qū)內(nèi)的計(jì)算收斂,不需要依賴所有分區(qū)的數(shù)據(jù),可以并行地在不同節(jié)點(diǎn)進(jìn)行計(jì)算。所以它的失敗回復(fù)也更有效,因?yàn)樗恍枰匦掠?jì)算丟失的parentpartition即可窄依賴介紹窄依賴對(duì)優(yōu)化的優(yōu)勢(shì)/021.依賴往往對(duì)應(yīng)著Shuffle操作,需要在運(yùn)行過(guò)程中將同一個(gè)父RDD的分區(qū)傳入到不同的子RDD分區(qū)中,中間可能涉及多個(gè)節(jié)點(diǎn)之間的數(shù)據(jù)傳輸;而窄依賴的每個(gè)父RDD的分區(qū)只會(huì)傳入到一個(gè)子RDD分區(qū)中,通常可以在一個(gè)節(jié)點(diǎn)內(nèi)完成轉(zhuǎn)換。2.當(dāng)RDD分區(qū)丟失時(shí)(某個(gè)節(jié)點(diǎn)故障),Spark會(huì)對(duì)數(shù)據(jù)進(jìn)行重算窄依賴對(duì)優(yōu)化的優(yōu)勢(shì)3.對(duì)于窄依賴,由于父RDD的一個(gè)分區(qū)只對(duì)應(yīng)一個(gè)子RDD分區(qū),這樣只需要重算和子RDD分區(qū)對(duì)應(yīng)的父RDD分區(qū)即可,所以這個(gè)重算對(duì)數(shù)據(jù)的利用率是100%的。4.對(duì)于依賴,重算的父RDD分區(qū)對(duì)應(yīng)多個(gè)子RDD分區(qū)的,這樣實(shí)際上父RDD中只有一部分的數(shù)據(jù)是被用于恢復(fù)這個(gè)丟失的子RDD分區(qū)的,另一部分對(duì)應(yīng)子RDD的其他未丟失分區(qū),這就造成了多余的計(jì)算;更一般的,寬依賴中子RDD分區(qū)通常來(lái)自多個(gè)父RDD分區(qū),極端情況下,所有的父RDD分區(qū)都要進(jìn)行重新計(jì)算。窄依賴對(duì)優(yōu)化的優(yōu)勢(shì)1.窄依賴介紹

2.窄依賴對(duì)優(yōu)化的優(yōu)勢(shì)

寬依賴寬依賴介紹寬依賴和窄依賴的區(qū)別寬依賴介紹/01

寬依賴是指父RDD的每個(gè)分區(qū)都可能被多個(gè)子RDD分區(qū)所使用,子RDD分區(qū)通常對(duì)應(yīng)所有的父RDD分區(qū),即一個(gè)父rdd的數(shù)據(jù)到多個(gè)子rdd(一對(duì)多)。

寬依賴介紹

寬依賴分為兩種情況:1個(gè)父RDD對(duì)應(yīng)非全部多個(gè)子RDD分區(qū),比如groupByKey,reduceByKey,sortByKey。1個(gè)父RDD對(duì)應(yīng)所有子RDD分區(qū),比如未經(jīng)協(xié)同劃分的join。寬依賴介紹

寬依賴需要所有的父分區(qū)都是可用的,必須等RDD的parentpartition數(shù)據(jù)全部ready之后才能開始計(jì)算,可能還需要調(diào)用類似MapReduce之類的操作進(jìn)行跨節(jié)點(diǎn)傳遞。從失敗恢復(fù)的角度看,shuffledependency牽涉RDD各級(jí)的多個(gè)parentpartition。寬依賴介紹寬依賴和窄依賴的區(qū)別/021、窄依賴是將分區(qū)聚合到一起,收攏數(shù)據(jù),這樣就可以考慮到一些算子做此功能比如:map,filter,union,join(父RDD是hash-partitioned),mapPartitions,mapValues;2、而寬依賴則不同,寬依賴將分區(qū)數(shù)據(jù)進(jìn)行打散分開,走shuffle機(jī)制與mapreduce相同。他主要將一些數(shù)據(jù)進(jìn)行洗牌和重新分組發(fā)牌。比如算子做此功能:groupByKey,join(父RDD不是hash-partitioned),partitionBy,sort寬依賴和窄依賴的區(qū)別1.寬依賴介紹

2.寬依賴和窄依賴的區(qū)別

RDD的依賴關(guān)系RDD依賴關(guān)系的本質(zhì)依賴關(guān)系下的數(shù)據(jù)流視圖RDD依賴關(guān)系的本質(zhì)/01

由于RDD是粗粒度的操作數(shù)據(jù)集,每個(gè)Transformation操作都會(huì)生成一個(gè)新的RDD,所以RDD之間就會(huì)形成類似流水線的前后依賴關(guān)系;RDD依賴關(guān)系的本質(zhì)

在spark中,RDD之間存在兩種類型的依賴關(guān)系:窄依賴和寬依賴;如圖所示顯示了RDD之間的依賴關(guān)系。RDD依賴關(guān)系的本質(zhì)依賴關(guān)系下的數(shù)據(jù)流視圖/02

如下圖是RDD依賴關(guān)系下的數(shù)據(jù)流視圖:依賴關(guān)系下的數(shù)據(jù)流視圖

在spark中,會(huì)根據(jù)RDD之間的依賴關(guān)系將DAG圖劃分為不同的階段,對(duì)于窄依賴,由于partition依賴關(guān)系的確定性,partition的轉(zhuǎn)換處理就可以在同一個(gè)線程里完成,窄依賴就被spark劃分到同一個(gè)stage中,而對(duì)于寬依賴,只能等父RDDshuffle處理完成后,下一個(gè)stage才能開始接下來(lái)的計(jì)算。依賴關(guān)系下的數(shù)據(jù)流視圖

因此spark劃分stage的整體思路是:從后往前推,遇到寬依賴就斷開,劃分為一個(gè)stage;遇到窄依賴就將這個(gè)RDD加入該stage中。依賴關(guān)系下的數(shù)據(jù)流視圖

在spark中,Task的類型分為2種:ShuffleMapTask和ResultTask;簡(jiǎn)單來(lái)說(shuō),DAG的最后一個(gè)階段會(huì)為每個(gè)結(jié)果的partition生成一個(gè)ResultTask,即每個(gè)Stage里面的Task的數(shù)量是由該Stage中最后一個(gè)RDD的Partition的數(shù)量所決定的,而其余所有階段都會(huì)生成ShuffleMapTask;之所以稱之為ShuffleMapTask,是因?yàn)樗枰獙⒆约旱挠?jì)算結(jié)果通過(guò)shuffle到下一個(gè)stage中;依賴關(guān)系下的數(shù)據(jù)流視圖1.RDD依賴關(guān)系的本質(zhì)

2.依賴關(guān)系下的數(shù)據(jù)流視圖

RDD的緩存RDD的緩存機(jī)制使用緩存RDD的緩存機(jī)制/01

默認(rèn)情況下,RDD只使用一次,用完即扔,再次使用時(shí)需要重新計(jì)算得到,而持久化操作避免了這里的重復(fù)計(jì)算,實(shí)際測(cè)試也顯示持久化對(duì)性能提升明顯。RDD的緩存機(jī)制

假設(shè)首先進(jìn)行了RDD0→RDD1→RDD2的計(jì)算作業(yè),那么計(jì)算結(jié)束時(shí),RDD1就已經(jīng)緩存在系統(tǒng)中了。在進(jìn)行RDD0→RDD1→RDD3的計(jì)算作業(yè)時(shí),由于RDD1已經(jīng)緩存在系統(tǒng)中,因此RDD0→RDD1的轉(zhuǎn)換不會(huì)重復(fù)進(jìn)行。RDD的緩存機(jī)制

計(jì)算作業(yè)只須進(jìn)行RDD1→RDD3的計(jì)算就可以了,因此計(jì)算速度可以得到很大提升。RDD的緩存機(jī)制使用緩存/02

使用如下代碼使用緩存:scala>importorg.apache.spark.storage._scala>valrdd1=sc.makeRDD(1to5)scala>rdd1.cache//cache只有一種默認(rèn)的緩存級(jí)別,即MEMORY_ONLYscala>rdd1.persist(StorageLevel.MEMORY_ONLY)使用緩存

Spark會(huì)自動(dòng)監(jiān)控每個(gè)節(jié)點(diǎn)上的緩存數(shù)據(jù),然后使用least-recently-used(LRU)機(jī)制來(lái)處理舊的緩存數(shù)據(jù)。如果想手動(dòng)清理這些緩存的RDD數(shù)據(jù)而不是去等待它們被自動(dòng)清理掉,使用RDD.unpersist()方法。使用緩存1.RDD緩存機(jī)制

2.使用緩存RDD的CheckPointCheckPoint介紹RDD的檢查點(diǎn)(checkpoint)機(jī)制CheckPoint介紹/01

檢查點(diǎn)是為了通過(guò)lineage(血統(tǒng))做容錯(cuò)的輔助,lineage過(guò)長(zhǎng)會(huì)造成容錯(cuò)成本過(guò)高,這樣就不如在中間階段做檢查點(diǎn)容錯(cuò),如果之后有節(jié)點(diǎn)出現(xiàn)問(wèn)題而丟失分區(qū),從做檢查點(diǎn)的RDD開始重做Lineage,就會(huì)減少開銷。

CheckPoint介紹

在設(shè)置檢查點(diǎn)之后,該RDD之前的有依賴關(guān)系的父RDD都會(huì)被銷毀,下次調(diào)用的時(shí)候直接從檢查點(diǎn)開始計(jì)算。

溫馨提示

  • 1. 本站所有資源如無(wú)特殊說(shuō)明,都需要本地電腦安裝OFFICE2007和PDF閱讀器。圖紙軟件為CAD,CAXA,PROE,UG,SolidWorks等.壓縮文件請(qǐng)下載最新的WinRAR軟件解壓。
  • 2. 本站的文檔不包含任何第三方提供的附件圖紙等,如果需要附件,請(qǐng)聯(lián)系上傳者。文件的所有權(quán)益歸上傳用戶所有。
  • 3. 本站RAR壓縮包中若帶圖紙,網(wǎng)頁(yè)內(nèi)容里面會(huì)有圖紙預(yù)覽,若沒(méi)有圖紙預(yù)覽就沒(méi)有圖紙。
  • 4. 未經(jīng)權(quán)益所有人同意不得將文件中的內(nèi)容挪作商業(yè)或盈利用途。
  • 5. 人人文庫(kù)網(wǎng)僅提供信息存儲(chǔ)空間,僅對(duì)用戶上傳內(nèi)容的表現(xiàn)方式做保護(hù)處理,對(duì)用戶上傳分享的文檔內(nèi)容本身不做任何修改或編輯,并不能對(duì)任何下載內(nèi)容負(fù)責(zé)。
  • 6. 下載文件中如有侵權(quán)或不適當(dāng)內(nèi)容,請(qǐng)與我們聯(lián)系,我們立即糾正。
  • 7. 本站不保證下載資源的準(zhǔn)確性、安全性和完整性, 同時(shí)也不承擔(dān)用戶因使用這些下載資源對(duì)自己和他人造成任何形式的傷害或損失。

最新文檔

評(píng)論

0/150

提交評(píng)論