一、創建RDD
兩種方式:
1.從檔案系統中加載資料創建RDD
Spark采用textFile()方法來從檔案系統中加載資料創建RDD,該方法把檔案的URI作為引數,這個URI可以是:
- 本地檔案系統的地址
- 或者是分布式檔案系統HDFS的地址
- 或者是Amazon S3的地址等等
2. 通過并行集合(串列)創建RDD
可以呼叫SparkContext的parallelize方法,在Driver中一個已經存在的集合 (串列)上創建,
舉個栗子:
第一種:
lines = sc.textFile("file:///usr/local/spark/mycode/rdd/word.txt") or sc.textFile("hdfs://localhost:9000/user/hadoop/word.txt")
>>> lines.foreach(print)
Hadoop is good
Spark is fast
Spark is bette第二種:
>>> array = [1,2,3,4,5]
>>> rdd = sc.parallelize(array)
>>> rdd.foreach(print)
二、轉換函式
1.filter()
.filter(func):篩選出滿足函式func的元素,并回傳一個新的資料集
>>>lines = sc.textFile("file:/l/usr/local/spark/mycode/rdd/word.txt")
>>>linesWithSpark = lines.filter(lambda line: "Spark" in line)
>>> linesWithSpark.foreach(print)
Spark is fast
Spark is better
2.map()
.map(func):將每個元素傳遞到函式func中,并將結果回傳為一個新的資料集,一個純粹的轉換操作
第一個栗子:
>>>data=[1,2,3,4,5]
>>> rdd1=sc.parallelize(data)
>>> rdd2=rdd1.map(lambda x:x+1)
>>> rdd2.foreach(print)第二個栗子:
>>>lines =sc.textFile("file:lllusr/local/spark/mycode/rdd/word.txt"')
>>>words = lines.map(lambda line:line.split(" "))
>>>words.foreach(print)
['Hadoop', ' is' , 'good']['Spark', 'is', ' fast]['Spark', 'is', ' better']
3.flatMap()
.flatMap(func):與map相似,但每個輸入元素都可以映射到0或多個輸出結果;先執行Map在執行flat拍扁其中的每個元素
第一個栗子:
>>>lines = sc.textFile("file:/llusr/localspark/mycode/rdd/word.txt")
>>>words = lines.flatMap(lambda line:line.split(" "))第二個栗子:
>>>words = sc.parallelize([("Hadoop",1),("is",1).("good",1),.... ("Spark",1).("is",1),("fast",1),("Spark",1),("is",1),("better",1)])
>>> words1 = words.groupByKey
>>> words1.foreach(print)
('Hadoop',<pyspark.resultiterable.ResultIterable object at 0x7fb210552c88>)
('better',<pyspark.resultiterable.ResultIterable object at 0x7fb210552e80>)
('fast',<pyspark.resultiterable.ResultIterable object at 0x7fb210552c88>)
('good',<pyspark.resultiterable.ResultIterable object at 0x7fb210552c88>)
('Spark',<pyspark.resultiterable.ResultIterable object at 0x7fb210552f98>)
('is',<pyspark.resultiterable.ResultIterable object at Ox7fb210552e10>)
4.reduceByKey()
.reduceByKey(func):應用于(K,V)鍵值對的資料集時,回傳一個新的(K,V)形式的資料集,其中每個值是將每個key傳遞到函式func中進行聚合后的結果
>>>words = sc.parallelize([("Hadoop",1),("is",1),("good",1),("Spark",1), !.... ("is",1),("fast"",1),("Spark",1),("is",1),("better",1)])
>>> words1 = words.reduceByKey(lambda a,b:a+b)
>>> words1.foreach(print)
('good', 1)
('Hadoop', 1)('better', 1)('Spark', 2)('fast', 1)('is',3)
tips:
- groupByKey也是對每個key進行操作,但只生成一個sequence, groupByKey本身不能自定義函式,需要先用groupByKey生成RDD,然后才能對此RDD通過map進行自定義函式操作.
- reduceByKey用于對每個key對應的多個value進行merge操作,最重要的是它能夠在本地先進行merge操作,并且merge操作可以通過函式自定義.
5.keys()
.keys():回傳鍵值
>>> list=[("Hadoop",1),("Spark",1),("Hive",1),("Spark",1)]
>>>pairRDD= sc.parallelize(list)
>>> pairRDD.keys().foreach(print)
Hadoop
Spark
Hive
Spark
6.values()
.values():回傳值
>>> list =[("Hadoop",1),("Spark",1),("Hive",1),("Spark",1)]
>>> pairRDD = sc.parallelize(list)
>>>pairRDD.values(.foreach(print)
1111
7.sortByKey()
.sortByKey():回傳一個根據鍵排序的RDD,默認升序排序,降序sortByKey(False)
第一個栗子:
>>>list=[("Hadoop",1).("Spark",1),("Hive",1),(""Spark",1)]
>>>pairRDD= sc.parallelize(list)
>>> pairRDD.foreach(print)
('Hadoop', 1)
('Spark', 1)
('Hive', 1)
('Spark', 1)第二個栗子:
>>>pairRDD.sortByKey.foreach(print)
('Hadoop', 1)
('Hive', 1)('Spark', 1)('Spark', 1)
>>>d1.reduceByKey(lambda a,b:a+b).sortByKey(False).collect()
[('g',21), ('f,29), ('e',17), ('d', 9), ('c',27),('b',38), ('a',42)]
8.sortBy()
.sortBy(func):按func自定義排序
>>>d1.reduceByKey(lambda a,b:a+b).sortBy(lambda x: x[0],False).collect(('g',21), ('f,29), ('e',17), ('d', 9), ('c',27), ('b',38), ('a', 42))
>>>d1.reduceByKey(lambda a,b:a+b).sortBy(lambda x: x[1],False).collect()
[('a',42), ('b',38), ('f,29), ('c',27),('g', 21), ('e',17), ('d', 9)]
9.mapValues()
.mapValues(func):對鍵值對RDD中的每個value都應用—個函式,key不會發生變化
>>> list=[("Hadoop",1),("Spark",1),("Hive",1),("Spark",1)]
>>>pairRDD= sc.parallelize(list)
>>>pairRDD1 = pairRDD.mapValues(lambda x;x+1)
>>>pairRDD1.foreach(print)
('Hadoop',2)
('Spark',2)('Hive',2)('Spark', 2)
10.join()
.join(): join就表示內連接,對于內連接,對于給定的兩個輸入資料集(K,V1)和(K,V2),只有在兩個資料集中都存在的key才會被輸出,最終得到一個:(K,(V1,V2))型別的資料集,
>>> pairRDD1 = sc. parallelize([("spark",1),("spark",2),("hadoop",3),("hadoop",5)])
>>>pairRDD2= sc.parallelize([("spark","fast")))
>>>pairRDD3 = pairRDD1.join(pairRDD2)
>>>pairRDD3.foreach(print)
('spark', (1, 'fast'))
('spark', (2, 'fast'))
11.distinct()
.distinct():去除重復值,一般用于在資料讀入時執行該操作
>>> RDD = sc. parallelize([("spark",1),("spark",1),("spark",2),("hadoop",3),("hadoop",5)]).map(lambda x:s.strip()).distinct()
>>>RDD.foreach(print)
("spark",1)
("spark",2)
("hadoop",3)
("hadoop",5)
二、常見的行動操作:
- count(回傳資料集中的元素個數
- collect0以陣列的形式回傳資料集中的所有元素
- first()回傳資料集中的第一個元素
- take(n)以陣列的形式回傳資料集中的前n個元素
- reduce(func)通過函式func(輸入兩個引數并回傳一個值)聚合資料集中的元素
- foreach(func)將資料集中的每個元素傳遞到函式func中運行
>>>rdd = sc.parallelize([1,2,3,4,5)
>>> rdd.countO
5
>>> rdd.first()
1
>>>rdd.take(3)
[1,2,3]
>>>rdd.reduce(lambda a,b:a+b)
15
>>> rdd.collect()
[1,2,3,4,5]
>>>rdd.foreach(lambda elem:print(elem))
12345
三、題外話:
A.持久化
在Spark中,RDD采用惰性求值的機制,每次遇到行動操作,都會從頭開始執 行計算,每次呼叫行動操作,都會觸發一次從頭開始的計算,這對于迭代計 算而言,代價是很大的,迭代計算經常需要多次重復使用同一組資料,通過持久化(快取)機制避免這種重復計算的開銷,
.persist():標記為持久化,在第一次行動操作時執行----->.unpersist():手動地把持久化的RDD從快取中移除
>>> list =["Hadoop","Spark","Hive"]
>>> rdd = sc.parallelize(list)
>>>rdd.cache) #會呼叫persist(MEMORY_ONLY),但是,陳述句執行到這里,并不會快取rdd,因為這時rdd還沒有被計算生成
>>> print(rdd.count() #第一次行動操作,觸發一次真正從頭到尾的計算,這時上面的rdd.cache()才會被執行,把這個rdd放到快取中
3
>>>print(','.join(rdd.collectO)) #第二次行動,不需要觸發從頭到尾的計算,只需要重復使用上面快取中的rdd
Hadoop,Spark,Hive
B.磁區
RDD是彈性分布式資料集,通常RDD很大,會被分成很多 個磁區,分別保存在不同的節點上,磁區的作用主要是:增加并行度;減少通信開銷,
磁區原則:RDD磁區的一個原則是使得磁區的個數盡量等于集群中的CPU核心 (core)數目
磁區個數:
(1)創建RDD時手動指定磁區個數
sc.textFile(path, partitionNum) 其中,path引數用于指定要加載的檔案的地址,partitionNum引數用于 指定磁區個數,例如:

(2)使用reparititon方法重新設定磁區個數
通過轉換操作得到新 RDD 時,直接呼叫 repartition 方法即可,例如:

轉載請註明出處,本文鏈接:https://www.uj5u.com/qita/382828.html
標籤:其他
