最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

Spark-Shell的啟動(dòng)與運(yùn)行實(shí)現(xiàn)過程

 更新時(shí)間:2026年01月17日 09:50:49   作者:LMY~~  
文章介紹了如何啟動(dòng)Spark,并演示了Spark RDD的基本操作,包括從文件系統(tǒng)和集合創(chuàng)建RDD,以及使用RDD編程API進(jìn)行數(shù)據(jù)處理,如map、filter、flatMap、reduceByKey等,通過實(shí)例,展示了如何對(duì)RDD進(jìn)行基本的數(shù)學(xué)運(yùn)算、字符串處理和數(shù)據(jù)統(tǒng)計(jì)

一、啟動(dòng)spark

1.先啟動(dòng)zookeeper

三臺(tái)虛擬機(jī)都要啟動(dòng)

zkServer.sh start

2.啟動(dòng)hadoop

start-all.sh

3.啟動(dòng)spark

在spark的根目錄下輸入

sbin/start-all.sh
spark-shell

二、Spark Rdd的簡(jiǎn)單操作

1.從文件系統(tǒng)加載數(shù)據(jù)創(chuàng)建ADD

(1)從Linux本地文件系統(tǒng)加載數(shù)據(jù)創(chuàng)建RDD——textFile(path)

val rdd = sc.textFile("file:///root/word.txt")
rdd.collect()   //查看命令

(2)從HDFS中加載數(shù)據(jù)創(chuàng)建RDD

val rdd = sc.textFile("/spark/test/word.txt")
rdd.collect()

scala> val rdd = sc.textFile("/spark/test/word.txt")
rdd: org.apache.spark.rdd.RDD[String] = /spark/test/word.txt MapPartitionsRDD[60] at textFile at :24
scala> rdd.collect()
res27: Array[String] = Array(hello java, hello hadoop, hello mysql)

2.通過集合創(chuàng)建RDD——prarallize()

從一個(gè)已經(jīng)存在的集合、數(shù)組,通過sarkContext對(duì)象調(diào)用parallelize的方法創(chuàng)建RDD,

val array =Array(1,2,3,4,5)
val arrRdd = sc.parallelize(array)
arrRdd.collect()

scala> val array =Array(1,2,3,4,5)
array: Array[Int] = Array(1, 2, 3, 4, 5)
scala> val arrRdd = sc.parallelize(array)
arrRdd: org.apache.spark.rdd.RDD[Int] = ParallelCollectionRDD[61] at parallelize at :26
scala> arrRdd.collect()
res29: Array[Int] = Array(1, 2, 3, 4, 5)

3.RDD的處理

一些RDD編程API

命令含義
map()返回一個(gè)新的rdd,由()轉(zhuǎn)換后組成
filter()過濾,由()函數(shù)計(jì)算后返回值為true的元素組成
flatMap()類似于map,但輸入元素可以被映射,用于詞頻拆分
union()相當(dāng)于數(shù)學(xué)中集合的并集
intersection()相當(dāng)于數(shù)學(xué)中集合的交集
distinct()去重操作后返回一個(gè)新的rdd
groupByKey()返回一個(gè)(l,iterator[數(shù)據(jù)類型]) 的rdd
reduceByKey()在一個(gè)(k,v)對(duì)的rdd上調(diào)用,返回一個(gè)新的(k,v)對(duì)rdd,用于詞頻統(tǒng)計(jì)
sortByKey在一個(gè)(k,v)對(duì)的rdd上調(diào)用,第二個(gè)值為true時(shí)按從小到大排序,false為從大到小排序
join()返回相同的key對(duì)應(yīng)的所有元素,如(K,(V,W))

(1)案例1

通過并進(jìn)行生成rdd
val rdd1 =List(5, 6, 4, 7, 3, 8, 2, 9, 1, 10)
val rdd2 =sc.parallelize(rdd1)

scala> val rdd1 =List(5,6,4,7,3,8,2,9,1,10)
rdd1: List[Int] = List(5, 6, 4, 7, 3, 8, 2, 9, 1, 10)
scala> val rdd2 =sc.parallelize(rdd1)
rdd2: org.apache.spark.rdd.RDD[Int] = ParallelCollectionRDD[0] at parallelize at :26
scala> rdd2.collect()
res3: Array[Int] = Array(5, 6, 4, 7, 3, 8, 2, 9, 1, 10)

對(duì)rdd1里的每一個(gè)元素乘2然后排序
val rdd3=rdd2.map(x=>x*2)
rdd3.collect()

scala> val rdd3=rdd2.map(x=>x*2)
rdd3: org.apache.spark.rdd.RDD[Int] = MapPartitionsRDD[1] at map at :28
scala> rdd3.collect()
res4: Array[Int] = Array(10, 12, 8, 14, 6, 16, 4, 18, 2, 20)

val rdd4=rdd3.sortBy(x=>x,true)
 rdd4.collect()

scala> val rdd4=rdd3.sortBy(x=>x,true)
rdd4: org.apache.spark.rdd.RDD[Int] = MapPartitionsRDD[4] at sortBy at :30
scala> rdd4.collect()
res5: Array[Int] = Array(2, 4, 6, 8, 10, 12, 14, 16, 18, 20)

(2)實(shí)例2

val rdd1 = sc.parallelize(Array("a b c", "d e f", "h i j"))
rdd1.collect()

scala> val rdd1 = sc.parallelize(Array(“a b c”, “d e f”, “h i j”))
rdd1: org.apache.spark.rdd.RDD[String] = ParallelCollectionRDD[5] at parallelize at :24
scala> rdd1.collect()
res6: Array[String] = Array(a b c, d e f, h i j)

將rdd1里面的每一個(gè)元素先切分在壓平
val rdd2 = rdd1.flatMap(x=>x.split(" "))
rdd2.collect

scala> val rdd2 = rdd1.flatMap(x=>x.split(" "))
rdd2: org.apache.spark.rdd.RDD[String] = MapPartitionsRDD[8] at flatMap at :26
scala> rdd2.collect()
res7: Array[String] = Array(a, b, c, d, e, f, h, i, j)

(3)實(shí)例3

計(jì)數(shù)word.txt中單詞的數(shù)量

 val rdd = sc.textFile("/spark/test/word.txt")   //從HDFS中加載數(shù)據(jù)創(chuàng)建RDD
 val rdd1=rdd.flatMap(x=>x.split(" "))			//將rdd用空格分開	
 val rdd2=rdd1.map(x=>(x,1))					//將不同的單詞(k,v)=(k,1)
 rdd2.collect()
 val rdd3=rdd2.groupByKey()						//相同的k放到一起
 rdd3.collect()
 val rdd4=rdd2.reduceByKey((a,b)=>a+b)			//將單詞進(jìn)行數(shù)量統(tǒng)計(jì)
 rdd4.collect()

scala> val rdd = sc.textFile("/spark/test/word.txt")
rdd: org.apache.spark.rdd.RDD[String] = /spark/test/word.txt MapPartitionsRDD[10] at textFile at :24
scala> val rdd1=rdd.flatMap(x=>x.split(" "))
rdd1: org.apache.spark.rdd.RDD[String] = MapPartitionsRDD[11] at flatMap at :26
scala> val rdd2=rdd1.map(x=>(x,1))
rdd2: org.apache.spark.rdd.RDD[(String, Int)] = MapPartitionsRDD[12] at map at :28
scala> rdd2.collect()
res9: Array[(String, Int)] = Array((hello,1), (java,1), (hello,1), (hadoop,1), (hello,1), (mysql,1))
scala> val rdd3=rdd2.groupByKey()
rdd3: org.apache.spark.rdd.RDD[(String, Iterable[Int])] = ShuffledRDD[13] at groupByKey at :30
scala> rdd3.collect()
res10: Array[(String, Iterable[Int])] = Array((hadoop,CompactBuffer(1)), (mysql,CompactBuffer(1)), (hello,CompactBuffer(1, 1, 1)), (java,CompactBuffer(1)))
scala> val rdd4=rdd2.reduceByKey((a,b)=>a+b)
rdd4: org.apache.spark.rdd.RDD[(String, Int)] = ShuffledRDD[14] at reduceByKey at :30
scala> rdd4.collect()
res11: Array[(String, Int)] = Array((hadoop,1), (mysql,1), (hello,3), (java,1))

(4)實(shí)例4

val rdd1 = sc.parallelize(List(("tom", 1), ("jerry", 3), ("kitty", 2)))
val rdd2 = sc.parallelize(List(("jerry", 2), ("tom", 1), ("shuke", 2)))
val rdd3 = rdd1.join(rdd2)		//求join
val rdd4 = rdd1.union(rdd2)		//求并集
val rdd5 = rdd4.groupByKey()	//按key分組
rdd5.collect					//查看

res11: Array[(String, Int)] = Array((hadoop,1), (mysql,1), (hello,3), (java,1))
scala> val rdd1 = sc.parallelize(List((“tom”, 1), (“jerry”, 3), (“kitty”, 2)))
rdd1: org.apache.spark.rdd.RDD[(String, Int)] = ParallelCollectionRDD[15] at parallelize at :24
scala> val rdd2 = sc.parallelize(List((“jerry”, 2), (“tom”, 1), (“shuke”, 2)))
rdd2: org.apache.spark.rdd.RDD[(String, Int)] = ParallelCollectionRDD[16] at parallelize at :24
scala> val rdd3 = rdd1.join(rdd2)
rdd3: org.apache.spark.rdd.RDD[(String, (Int, Int))] = MapPartitionsRDD[19] at join at :28
scala> val rdd4 = rdd1.union(rdd2)
rdd4: org.apache.spark.rdd.RDD[(String, Int)] = UnionRDD[20] at union at :28
scala> rdd3.collect()
res12: Array[(String, (Int, Int))] = Array((tom,(1,1)), (jerry,(3,2)))
scala> rdd4.collect()
res13: Array[(String, Int)] = Array((tom,1), (jerry,3), (kitty,2), (jerry,2), (tom,1), (shuke,2))
scala> val rdd5 = rdd4.groupByKey()
rdd5: org.apache.spark.rdd.RDD[(String, Iterable[Int])] = ShuffledRDD[21] at groupByKey at :30
scala> rdd5.collect()
res14: Array[(String, Iterable[Int])] = Array((tom,CompactBuffer(1, 1)), (jerry,CompactBuffer(3, 2)), (shuke,CompactBuffer(2)), (kitty,CompactBuffer(2)))

(5)實(shí)例5

val rdd1 = sc.parallelize(List(5, 6, 4, 3))
val rdd2 = sc.parallelize(List(1, 2, 3, 4))
val rdd3 = rdd1.union(rdd2)			//求并集
rdd3.collect()
val rdd4 = rdd1.intersection(rdd2) //求交集
rdd4.collect()
val rdd5 = rdd4.distinct()			//去重
rdd5.collect						//查看

scala> val rdd1 = sc.parallelize(List(5, 6, 4, 3))
rdd1: org.apache.spark.rdd.RDD[Int] = ParallelCollectionRDD[22] at parallelize at :24
scala> val rdd2 = sc.parallelize(List(1, 2, 3, 4))
rdd2: org.apache.spark.rdd.RDD[Int] = ParallelCollectionRDD[23] at parallelize at :24
scala> val rdd3 = rdd1.union(rdd2)
rdd3: org.apache.spark.rdd.RDD[Int] = UnionRDD[24] at union at :28
scala> rdd3.collect()
res15: Array[Int] = Array(5, 6, 4, 3, 1, 2, 3, 4)
scala> val rdd4 = rdd1.intersection(rdd2)
rdd4: org.apache.spark.rdd.RDD[Int] = MapPartitionsRDD[30] at intersection at :28
scala> rdd4.collect()
res16: Array[Int] = Array(4, 3)
scala> val rdd5 = rdd4.distinct()
rdd5: org.apache.spark.rdd.RDD[Int] = MapPartitionsRDD[33] at distinct at :30
scala> rdd5.collect()
res17: Array[Int] = Array(4, 3)

總結(jié)

以上為個(gè)人經(jīng)驗(yàn),希望能給大家一個(gè)參考,也希望大家多多支持腳本之家。

相關(guān)文章

  • Linux命令dos2unix命令示例詳解(將DOS格式文本文件轉(zhuǎn)換成Unix格式)

    Linux命令dos2unix命令示例詳解(將DOS格式文本文件轉(zhuǎn)換成Unix格式)

    dos2unix命令 用來(lái)將DOS格式的文本文件轉(zhuǎn)換成UNIX格式的,而Unix格式的文本文件在Windows下用Notepad打開時(shí)會(huì)拼在一起顯示,本文介紹Linux命令dos2unix命令示例詳解(將DOS格式文本文件轉(zhuǎn)換成Unix格式),感興趣的朋友一起看看吧
    2024-04-04
  • Shell過濾器的具體使用

    Shell過濾器的具體使用

    這篇文章主要介紹了Shell過濾器的具體使用,文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來(lái)一起學(xué)習(xí)學(xué)習(xí)吧
    2021-03-03
  • 基于shell的if和else詳解

    基于shell的if和else詳解

    下面小編就為大家?guī)?lái)一篇基于shell的if和else詳解。小編覺得挺不錯(cuò)的,現(xiàn)在就分享給大家,也給大家做個(gè)參考。一起跟隨小編過來(lái)看看吧
    2017-06-06
  • Shell腳本中獲取本機(jī)ip地址的3個(gè)方法

    Shell腳本中獲取本機(jī)ip地址的3個(gè)方法

    這篇文章主要介紹了Shell腳本中獲取本機(jī)ip地址的3個(gè)方法,本文直接給出實(shí)現(xiàn)代碼,需要的朋友可以參考下
    2014-10-10
  • 用shell命令讀取與輸出數(shù)據(jù)的代碼

    用shell命令讀取與輸出數(shù)據(jù)的代碼

    本文為大家介紹使用shell命令進(jìn)行讀取與輸出數(shù)據(jù)的方法,其中涉及了文件輸出、重定向、管道等相關(guān)知識(shí),有興趣的朋友可以參考下
    2013-02-02
  • Linux下shell基本命令之grep用法及示例小結(jié)

    Linux下shell基本命令之grep用法及示例小結(jié)

    grep是Unix/Linux系統(tǒng)中用于文本搜索的強(qiáng)大工具,它可以忽略大小寫、顯示行號(hào)、反向選擇、遞歸搜索目錄等,本文就來(lái)介紹一下,感興趣的可以了解一下
    2024-12-12
  • shell腳本運(yùn)行5秒后自動(dòng)退出的代碼

    shell腳本運(yùn)行5秒后自動(dòng)退出的代碼

    shell腳本運(yùn)行5秒自動(dòng)退出的代碼,供大家學(xué)習(xí)參考
    2013-02-02
  • nvidia-smi命令詳解和一些高階技巧講解

    nvidia-smi命令詳解和一些高階技巧講解

    一般情況下用的比較多的就是nvidia-smi的命令,其實(shí)掌握了這一個(gè)命令也就能夠覆蓋絕大多數(shù)場(chǎng)景了,但是本質(zhì)求真務(wù)實(shí)的態(tài)度,本文調(diào)研了相關(guān)資料,整理了一些比較常用的nvidia-smi命令的其他用法,感興趣的朋友跟隨小編一起看看吧
    2023-01-01
  • linux?shell字符串截取的詳細(xì)總結(jié)(實(shí)用!)

    linux?shell字符串截取的詳細(xì)總結(jié)(實(shí)用!)

    在開發(fā)的時(shí)候經(jīng)常會(huì)自行寫一些小的腳本,其中就用到截取字符串的操作,這篇文章主要給大家介紹了關(guān)于linux?shell字符串截取的詳細(xì)方法,文中通過實(shí)例代碼介紹的非常詳細(xì),需要的朋友可以參考下
    2022-07-07
  • 詳解Linux定時(shí)任務(wù)Crontab的介紹與使用

    詳解Linux定時(shí)任務(wù)Crontab的介紹與使用

    linux內(nèi)置的cron進(jìn)程能幫我們實(shí)現(xiàn)這些需求,cron搭配shell腳本,非常復(fù)雜的指令也沒有問題。本文主要介紹了定時(shí)任務(wù)Crontab的使用,需要的可以學(xué)習(xí)一下
    2022-10-10

最新評(píng)論

杨浦区| 棋牌| 安庆市| 峨眉山市| 阜阳市| 龙胜| 合肥市| 武汉市| 泗洪县| 盐边县| 璧山县| 石狮市| 铁力市| 梅州市| 海南省| 宁化县| 高雄县| 宝鸡市| 五常市| 麻栗坡县| 乌鲁木齐县| 津市市| 贡山| 通化县| 晋城| 巴林右旗| 琼中| 德阳市| 咸丰县| 同江市| 萨迦县| 栖霞市| 渑池县| 巴楚县| 娱乐| 洛川县| 泗洪县| 北流市| 茌平县| 金昌市| 南靖县|