顯示具有 技術-BigData-Spark 標籤的文章。 顯示所有文章
顯示具有 技術-BigData-Spark 標籤的文章。 顯示所有文章

2020年4月20日 星期一

你知道在 Azure 上有幾種 On Demand 啟動 Spark 的方法嗎?



最近需要開始分析一些 Log ,最直覺的方式就是使用最熟悉的 Spark 來分析,於是開始研究最近有什麼方便在 Azure 啟動 Spark 的方式,在 AWS 和 GCP 上,之前就已經有研究過專門支援的 PaaS 服務:
我知道 Azure 有 HDInsight ,但是之前使用覺得沒有 GCP Dataproc 好用,不知道 2020 年的今天,有沒有什麼新的 Solution 呢?畢竟 Azure 的 "強項" 就是透過大量跟 3rd-party ISV 整合來壯大自己的服務

2019年11月11日 星期一

Apache Spark 3.0 Release Preview






一轉眼 Apache Spark 3.0-preview 已經推出了,眼看又快要跟時代脫節了...🙈
這次 3.0 的更新多到爆炸,所以再研究有什麼令人興奮的更新前,我們先來複習一下之前在 Spark+AI Summit 2019 Keynote 的預告有什麼需要我們注意的。

  • 更多向量矩陣運算和GPU支援
  • 跟 K8s 更緊密的整合
  • Koalas - panda dataframe

2019年4月25日 星期四

Spark+AI Summit 2019 Keynote 重點搶鮮看



熱騰騰的 Spark+AI Summit 2019 的影片陸續出爐了,讓我們先來看看 Keynote 的重點內容  - Reynold Xin (Databricks), Brooke Wenig (Databricks)


第一個重點就是針對Unify Data 處理和AI Databricks 做了什麼努力,去年他們提出 Hydrogen 為了讓Spark 能更方便跟各種 ML lib 串接。


而今年的 Spark 3.0 放了更多重點在於讓 jvm base 的底層可以支援更多向量矩陣運算和GPU支援。








再來就是因應Kubernetes 的崛起,Spark 勢必得更加密切的與Kubernetes整合。



再來就是 Spark 3.x 想要解決 Data scientist 的痛,因為 Data scientist 通常用 python + panda 在他們的個人電腦上建模型和測試,但是一旦要scale 就得重寫code porting 到 spark,此外雖然看起來都是dataframe 但是實際上理念卻是差很多,所以stackoverflow 上常常都是這些型態轉換的問題。




於是Databricks推出 Koalas: Panda DataFrame API on Spark,最神奇的就是只要把panda的任function 換個名稱koalas 就無痛轉移了....XD




相信Data scientist 和 Data engineer 一定很期待,也可以少很多工~










2019年1月7日 星期一

Google Cloud Dataproc 如何建立 Custom Image 加快 PySpark 部署環境速度



這陣子最常使用的GCP服務就是 Cloud Dataproc , Cloud Dataproc 是為了簡化Spark及Hadoop服務而設計,能讓使用者進行批次處理、查詢、資料串流及機器學習等工作,其自動化工具可協助使用者更快新增及更容易管理資料叢集,並且能在不使用時關閉,以降低成本,使企業能把心力花在資料分析的核心工作上。

對我來說他有幾個優點:
  • 不用自己維運一組Yarn Spark Cluster,更不用煩惱需要擴充配置的問題
  • 需要用就直接開,開完就砍掉,一切自動化又省錢

2018年6月16日 星期六

[Spark 學習小筆記] 什麼是 Hyperparameter Tunning? 有什麼方法?



不論在學習或是工作常常都經歷過一段知其然不知其所以然的階段,是否能突破其實就是看有沒有這個機運和決心去突破,話說之前在翻譯 Spark ML 那本書時,看到 Grid Search (網格搜尋) 和 Random Search (隨機搜尋) 的時候其實是一頭霧水,只知道是用來搜尋Hyperparameter 的演算法,但是原理和如何使用卻是一無所知,直到最近開始利用Spark 開發Machine Learning 專案,就慢慢開始有感覺了。

2018年5月10日 星期四

Spark 小技巧系列 - 讀取檔案過大發生 java.lang.NegativeArraySizeException 該怎麼處理?


雖然我們知道單一檔案不要太大,或太小,但是有時候人在江湖身不由己,如果遇到單一檔案太大時,系統可能就會噴出以下錯誤:


[WARN] BlockManager   : Putting block rdd_12_0 failed due to exception java.lang.NegativeArraySizeException.
[WARN] BlockManager   : Block rdd_12_0 could not be removed as it was not found on disk or in memory
[ERROR] Executor          : Exception in task 0.0 in stage 3.0 (TID 259)
java.lang.NegativeArraySizeException
        at org.apache.spark.unsafe.types.UTF8String.concatWs(UTF8String.java:901)
        at org.apache.spark.sql.catalyst.expressions.GeneratedClass$SpecificUnsafeProjection.apply(Unknown Source)
        at org.apache.spark.sql.execution.aggregate.AggregationIterator$$anonfun$generateResultProjection$1.apply(AggregationIterator.scala:234)
        at org.apache.spark.sql.execution.aggregate.AggregationIterator$$anonfun$generateResultProjection$1.apply(AggregationIterator.scala:223)
        at org.apache.spark.sql.execution.aggregate.ObjectAggregationIterator.next(ObjectAggregationIterator.scala:86)
        at org.apache.spark.sql.execution.aggregate.ObjectAggregationIterator.next(ObjectAggregationIterator.scala:33)
        at org.apache.spark.sql.catalyst.expressions.GeneratedClass$GeneratedIteratorForCodegenStage3.processNext(Unknown Source)
        at org.apache.spark.sql.execution.BufferedRowIterator.hasNext(BufferedRowIterator.java:43)
        at org.apache.spark.sql.execution.WholeStageCodegenExec$$anonfun$10$$anon$1.hasNext(WholeStageCodegenExec.scala:614)
        at org.apache.spark.sql.execution.columnar.InMemoryRelation$$anonfun$1$$anon$1.hasNext(InMemoryRelation.scala:139)
        at org.apache.spark.storage.memory.MemoryStore.putIteratorAsBytes(MemoryStore.scala:378)
        at org.apache.spark.storage.BlockManager$$anonfun$doPutIterator$1.apply(BlockManager.scala:1109)
        at org.apache.spark.storage.BlockManager$$anonfun$doPutIterator$1.apply(BlockManager.scala:1083)
        at org.apache.spark.storage.BlockManager.doPut(BlockManager.scala:1018)
        at org.apache.spark.storage.BlockManager.doPutIterator(BlockManager.scala:1083)
        at org.apache.spark.storage.BlockManager.getOrElseUpdate(BlockManager.scala:809)
        at org.apache.spark.rdd.RDD.getOrCompute(RDD.scala:335)
        at org.apache.spark.rdd.RDD.iterator(RDD.scala:286)
        at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:38)
        at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:324)
        at org.apache.spark.rdd.RDD.iterator(RDD.scala:288)
        at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:38)
        at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:324)
        at org.apache.spark.rdd.RDD.iterator(RDD.scala:288)
        at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:38)
        at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:324)
        at org.apache.spark.rdd.RDD.iterator(RDD.scala:288)
        at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:38)
        at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:324)
        at org.apache.spark.rdd.RDD.iterator(RDD.scala:288)
        at org.apache.spark.scheduler.ShuffleMapTask.runTask(ShuffleMapTask.scala:96)
        at org.apache.spark.scheduler.ShuffleMapTask.runTask(ShuffleMapTask.scala:53)
        at org.apache.spark.scheduler.Task.run(Task.scala:109)
        at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:345)
        at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
        at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
        at java.lang.Thread.run(Thread.java:748) 

2018年5月9日 星期三

Spark 小技巧系列 - left join 後把null 改為0


如果使用Spark 的 left outer join 遇到沒有的資料通常就會以NULL顯示,如下圖所示:


這時候如果我想要計算CTR = click/ impression 會發生什麼事?直接噴錯給你看,然後也不知道發生什麼事....

org.apache.spark.sql.AnalysisException: Resolved attribute(s) 'impressionCount,'clickCount missing from viewCount#965L,(impressionCount / total)#582,dsType#1189,rid#14,impressionCount#534L,recommendCount#1198L,clickCount#1441L,siteId#16 in operator 'Project [siteId#16, rid#14, impressionCount#534L, (impressionCount / total)#582, viewCount#965L, dsType#1189, recommendCount#1198L, clickCount#1441L, ('clickCount / 'impressionCount) AS CTR#2022]. Attribute(s) with the same name appear in the operation: impressionCount,clickCount. Please check if the right attribute(s) are used.;;
'Project [siteId#16, rid#14, impressionCount#534L, (impressionCount / total)#582, viewCount#965L, dsType#1189, recommendCount#1198L, clickCount#1441L, ('clickCount / 'impressionCount) AS CTR#2022]
+- AnalysisBarrier
      +- LogicalRDD [siteId#16, rid#14, impressionCount#534L, (impressionCount / total)#582, viewCount#965L, dsType#1189, recommendCount#1198L, clickCount#1441L], false


    at org.apache.spark.sql.catalyst.analysis.CheckAnalysis$class.failAnalysis(CheckAnalysis.scala:41)
    at org.apache.spark.sql.catalyst.analysis.Analyzer.failAnalysis(Analyzer.scala:91)
    at org.apache.spark.sql.catalyst.analysis.CheckAnalysis$$anonfun$checkAnalysis$1.apply(CheckAnalysis.scala:289)
    at org.apache.spark.sql.catalyst.analysis.CheckAnalysis$$anonfun$checkAnalysis$1.apply(CheckAnalysis.scala:80)
    at org.apache.spark.sql.catalyst.trees.TreeNode.foreachUp(TreeNode.scala:127)
    at org.apache.spark.sql.catalyst.analysis.CheckAnalysis$class.checkAnalysis(CheckAnalysis.scala:80)
    at org.apache.spark.sql.catalyst.analysis.Analyzer.checkAnalysis(Analyzer.scala:91)
    at org.apache.spark.sql.catalyst.analysis.Analyzer.executeAndCheck(Analyzer.scala:104)
    at org.apache.spark.sql.execution.QueryExecution.analyzed$lzycompute(QueryExecution.scala:57)
    at org.apache.spark.sql.execution.QueryExecution.analyzed(QueryExecution.scala:55)
    at org.apache.spark.sql.execution.QueryExecution.assertAnalyzed(QueryExecution.scala:47)
    at org.apache.spark.sql.Dataset$.ofRows(Dataset.scala:74)
    at org.apache.spark.sql.Dataset.org$apache$spark$sql$Dataset$$withPlan(Dataset.scala:3295)
    at org.apache.spark.sql.Dataset.select(Dataset.scala:1307)
    at org.apache.spark.sql.Dataset.withColumns(Dataset.scala:2192)
    at org.apache.spark.sql.Dataset.withColumn(Dataset.scala:2159)


其實原因就是有Null的存在,這時候只要使用以下技巧補零就可以了。

Dataset join1 = impression.join(broadcast(view), col, LeftOuter.toString())
                               .na()
                               .fill(0, new String[] {"viewCount"});

這段的意思就是會把null 的值,補上任何你想要的值,然後就解決了~


2018年5月3日 星期四

加速 Spark 寫到雲端儲存空間(AWS S3/ Azure blob/ Google GCS )的速度



不知道大家有沒有遇到過明明 Spark 程式都結束了,檔案也寫完了,Driver program 確還不會停止,感覺就當在那邊的經驗?

很多技術細節沒有遇到還真得不知道會有這種設計和改善的方法,會發現這個密技是因為下面幾個條件同時成立才注意到的:

1. 使用雲端儲存空間


為了節省成本,我們並沒有架設自己的HDFS Server,取而代之的都是把要分析的資料和結果儲存在雲端儲存空間( AWS S3/ Azure blob/ Google GCS)。這個選擇會大大降低維運成本和提高檔案的保存可靠度,不過缺點就是失去data locality 的速度優勢,而且每次分析都從雲端拉檔案下來也會花不少時間,所以就是以時間效能換取成本和可靠度。

2018年4月27日 星期五

[Spark 學習小筆記] 如何使用Java 實作 vectorization (1)


Vectorization 的重要性


還記得上Anfrew Ng 的deep-learning課程前幾堂課就講到vectorization的重要性,以及對於效能會有怎樣的影響,對於矩陣運算,最直覺的反應就是寫個for loop,然後針對每個emlement 去做運算,但是其實CPU&GPU 有專門的指令集可以用來平行化專門處理這種運算,於是透過 vectorization 就可以得到顯著的效能改善,下圖就是老師上課時用python 的範例,相差了400多倍!?

2018年4月24日 星期二

Apache Spark 學習三部曲:學會他,除錯他,寫好他



最近密集的在寫Spark 程式,感覺到終於該開始往下個階段邁進了,其實就像學習任何程式語言和Framework,Spark 學習也分三個步驟:
  1. 如何寫
  2. 如何調教/錯誤排除
  3. 如何寫的好

學會如何寫,網路上有不少的範例,不過大多是Scala和Python,如果要翻成Java版還需要額外花點功夫,等到開始寫一些程式丟到spark 上面跑,又會開始遇到一堆奇奇怪怪的錯誤訊息,比如說:
  • Futures timed out after [300 seconds]
  • This timeout is controlled by spark.executor.heartbeatInterval
  • cannot assign instance of java.lang.invoke.SerializedLambda to field org.apache.spark.sql.UDFRegistration$$anonfun$27.f
  • Initial job has not accepted any resources; check your cluster UI to ensure that workers are registered and have sufficient memory


這時候就要開始學習怎麼調教系統參數和錯誤排除,下面這個網站整理的蠻不錯的,把Spark performance tuning 分成幾大塊,分別是:
  • Data Serialization
  • Memory Tuning
  • Memory Managemnt
  • Data Structure Tuning
  • Garbage Collection Tuning




最後就是程式都會動了,但是怎麼要跑那麼久?這時候就應該開始瞭解Spark 許多底層的運作原理,要怎麼寫才是正確的,怎樣寫才會有比較好的效能,下面收一些在學習spark 時收集到的不錯資訊:



其它網路文章:


最後要提醒自己:

先求有,再求好,等到程式會動了,確定結果是大家想要的,有符合商業價值了,再來調教也不遲!

2018年3月13日 星期二

透過Java開發 Spark 2.x ML 的 LDA (Latent Dirichlet allocation) model 的感想

借用 Deep Learning 那門課的一張投影片來代表我最近在做的事,一直不斷的在轉圈圈,但是其實更多時候腦袋都是下圖這種轉圈圈:


主要的挑戰如下:
  1. 第一次用Spark 寫 ML 相關的程式,然後網路上的範例和討論幾乎都是用 scala 和python 寫的,為了要轉成Java版花了不少力氣。
  2. 資料散亂且難以理解,關於LDA 的文章幾乎都是論文等級了,充滿了難以下嚥的數學公式,往往都是直接略過,而且關於LDA 和 Spark ML 相關的中文討論又是以大陸居多,不過也因此挖到不少寶。

Spark ML 內建的Pipeline 主要分為以下四個步驟:



不過後來才發現Spark ML 的 pipeline 是給 supervisor learning 使用的,因為有 label 的 data  可以用來驗證訓練的結果是好還是不好,但是我們這次用到的不論是 LDA 或是 word2vec 都是屬於 unsupervised learning 沒有一個明確的基準點可以驗證,因次不適合用 cross validation 來找出最佳的 hyper parameter 。



那LDA 可以調的參數有哪些些呢?如下圖所示,主要是 K (topic number) , max iteration,doc Concentration和 topic Concentration,不過上網看了許多論文和討論,似乎影響最大的還是K 值。




原本以為只要把這些參數排列組合,找出最大 likelihood 或最小的 perplexity 結果居然看到以下的討論......



在研究中產生許多問題,順便把問題整理在這裡:

1. LDA 跟 word2vec 的差異在哪裡?

LDA 注重的是文章與文章間所有詞的關係,而word2vec 是詞與某篇文章上下文之間的關係,也就是說 word2vec 並沒有考慮語法層面的訊息,一篇文章被看成文字序列(word sequences),只考慮詞與詞之間的位置與邊界關係。

網路上找了一個有趣的解釋,假如你輸入HTC:

普通的word2vec 會找到:Android,cellphone,Taiwan,Google...
但是考慮語法後:Moto,Apple,Nokia,Xiaomi...

雖然看起來都有關連,但是本質上卻是不一樣,真得好難啊....

更多網路資訊:

關於LDA 與 Spark 相關文章:


 





2016年11月30日 星期三

翻譯的第一本書問世了 - Spark大數據分析新利器




這個blog 似乎很久沒有文章了,因為....忙著生小孩(無誤),此外今年還給自己挖了個大坑,翻譯一本技術書!!  (偷懶還那麼多藉口...)

看到自己翻譯的書終於上市了,心中真是五味雜陳啊.....
首先還是得感謝神一般的隊友 - Study 大大,如果沒有他這本書一定是翻譯不完也翻譯不好的XD
平日上班就很忙了,下班還要照顧小孩,真得很難抽出時間翻譯,熬夜翻譯的結果就是品質很糟,又得花更多的時間校稿修正,因此就得感謝松崗的編輯對我們的耐心。
最後當然要感謝自己的家人的包容,因為下班回家後和假日都在翻譯,無法全心幫忙照顧小孩..😂

話說回來只不過翻譯一本書,怎麼寫的好像自己是原作者一樣...😂
自己走過一遭有了以下感想:
  1. 寫書的作者真得很辛苦~很辛苦...QAQ
  2. 在台灣翻譯書的譯者真得都是佛心來著,因為翻譯真不是人幹的,費時費力收入又真得很有限XDrz..

總之還請大家多多支持購買~~(你們買再多我也沒有分紅啦XDrz..)
如果有任何翻譯的錯誤,或是字句不通順,也還請多多包涵~~Orz...

2013年10月2日 星期三

Apache Spark 0.8.0 正式發佈了!!



千呼萬喚始出來,Apache Spark 0.8.0 終於正式Release 了,這是是捐給ASF後的第一版Release ,也是有著重大更新的一版Release,讓我們來看看Release note 裡面有提到哪些新Feature。

1. Monitoring UI and Metrics (上圖便是新的介面)

2. Machine Learning Library (有取代Mahout 的味道唷~)

3.Python Improvements (就不用再用大陸的山寨版了~XD)

4. Hadoop YARN support

5. Revamped Job Scheduler

6. Easier Deployment and Linking

7. Expanded EC2 Capabilities

8. Improved Documentation

9. Other- Hadoop save functions now support an optional compression codec.

第九點倒是對我來說蠻重要的,如果增加了compression codec support,那就代表可能也可以support Encryption摟~:D


2013年8月25日 星期日

第一次玩Spark Shark 就上手 - 不負責任效能測試



既然都安裝好了,總是要來比較一下Shark效能,是不是真的如傳說中那麼威~~


孬孬免責聲明:此篇測試不是在很嚴謹的環境,也沒有Fine tune的狀況下做出簡單的測試比較,純粹提供參考,有興趣的人建議還是自行測試~:P

測試的環境:

機器:Dell Power-edge 的機器上開4台VM (每台設定4 core CPU 4G Ram)
環境:
  • CentOS6.4
  • Hadoop CDH4.1.x  (Hive 0.9.0)
  • Spark stand-alone mode

當一切安裝就緒就可以在Master 的UI上看到以下資訊:





測試案例 - 統計銀行用戶年均存款餘額的分佈 


  • 年均存款餘額:從當年度1/1 到結算日每天的存款餘額加起來除以365  (如果當天沒有餘額變更記錄,則以上一次變更餘額為本日餘額)
  • 統計分布,分別以下級距來統計客戶數量:0~10,000、10,000~100,000、100,000~1,000,000、1,000,000~10,000,000

下表欄位意義說明:
  • Record per day 代表一天會有幾筆用戶存款資料變更
  • Days 代表產生幾天份資料
  • Data Size 代表實際產生的 File size
時間則是產出統計結果所需要花的時間









Dataset
Record
pre day
100
1,000
10,000
50,000
100,000
Days
365
365
365
365
365
Data Size
9.2Mb
95Mb
926Mb
4.6G
9.2G
Hive
Sec
73.298
143
700
1184
Dead!
Shark
Sec
37
108
216
2747
Dead!


在一開始資料量小的時候,的確Shark 都比Hive快很多,但是隨著資料量變大,vm的記憶體被吃光光,開始吃到swap時,Shark 的效能就會往下掉,然後我最後側到9.2G的檔案得時候,vm就全部死光光了....Orz...


之後會在想辦法找實體機器(要擁有足夠的記憶體)來測試可能會比較準,另外如果加入YARN或是Mesos可能又會有不同的結果...

而且玩到這裡覺得越來越有趣了,也產生了更多問題需要搞清楚,比如說:

1. worker 之間有無溝通?溝通內容?
2. 詳細了解mesos 的task 如分配工作 (順便了解YARN)
3. 了解Spark 如何切割工作?
4. coarse-grained 是否可以開一個以上的work ?
 

且讓我們繼續看下去~

2013年8月12日 星期一

第一次玩Spark Shark 就上手 - Spark安裝篇




Shark 和 Spark 在安裝上充滿了彈性,有很多種組合方式,下圖就是Spark 和 Shark 可以搭配的安裝模式,與種類。

圖片來源:自行整理


Spark 如上圖所示,在運行上有好幾種模式:

1. 分散式模式:


1-1. 架構在Mesos上


而架構在Mesos 上又有分兩種模式:

1-1-1. fine-grained (Default)

根據官網的解釋,在Fine-grained模式下,每一個Spark task 就等於是一個Mesos Task ,所以可以在同一台機器上跑好幾個Spark Task,好處是方便動態調配資源,需要的時候再去啟動一個Task,缺點是啟動每一個Spark Task 會需要額外的資源花費,所以比較不適合需要low-latency 的application(像是互動式query 或是 web request)。


圖片來源:自行整理

1-1-2. coarse-grained

而跑在coarse-grained 模式下,每一台Mesos 所管理的機器上,只會啟動一組Spark Task,整體資源規劃是透過mesos 動態排程(dynamically schedule) 所屬的作小工作"mini-tasks",優點就是啟動每個工作所需的時間較少,但是需要這個Task必須事先就預留起來。 (這應該比較像Hadoop 事先就先規劃好每台機器有幾個 Job Tracker一樣?)



圖片來源:自行整理

1-2. 架構在YARN 上 (實驗性質)


架構再YARN上號稱跟架構再Mesos 上一樣簡單,但是我就沒有特別去研究的,有興趣的可以自行研究。


2. 單機模式


既然是第一次玩Spark 就上手,所以這篇文章會著重在第二種模式,也就是單機模式, 顧名思義不需要架構在Mesos或是YARN 的Cluster 上,也可以單獨運作於Hadoop 之外,單機就可以直接執行運算,在這個模式下就可以直接存取Local Disk 或是Hdfs(透過lib 去存取Hdfs),甚至是S3。不過要澄清一下,所謂的單機模式不代表只能跑一台,他同樣也是可以跑很多台變成Cluster 的形式,唯一的差別就是所有的slave worker 都是由Spark Master所控制(類似coarse-grained mode 每一台機器都預先裝好一組Spark worker)。

所有的Node 都是Master Node 透過ssh在控制,如下圖所示:


圖片來源:自行整理


安裝步驟



Pre-Requirement:

1. 安裝Java (這應該大家都會就跳過了)

2. 安裝Scala  (注意:目前Spark 0.7.3 版限定只能用scala-2.9.3)

# wget http://www.scala-lang.org/files/archive/scala-2.9.3.tgz
# tar xvf scala-2.9.3.tgz
# sudo mv scala-2.9.3 /usr/lib
# sudo ln -s /usr/lib/scala-2.9.3 /usr/lib/scala

設定Path 和 Scala home 編輯 /etc/profile.d/scala.sh

export SCALA_HOME=/usr/lib/scala
export PATH=$PATH:$SCALA_HOME/bin


3. 下載並安裝Spark

# wget http://spark-project.org/download/spark-0.7.3-prebuilt-cdh4.tgz
# tar zxvf spark-0.7.3-prebuilt-cdh4.tgz
# mv spark-0.7.3 /usr/lib/
# ln -s /usr/lib/spark-0.7.3 /usr/lib/spark

編輯~/.bashrc 加入spark home and path

export SPARK_HOME=/usr/lib/scala
export PATH=$PATH:$SPARK_HOME/bin

4.  設定成 Standalone Cluster 模式 (其他台機器也依照前面的三個步驟安裝)

我邊準備了三台VM要用來跑Cluster,分別是lab-hadoop-m1,lab-hadoop-m2,lab-hadoop-m3,預計讓m1 跑Spark Master,由於Spark 會透過ssh 去控制其他台機器,所以建議先設定讓lab-hadoop-m1這台機器不用輸入密碼,改使用private key 的方式登入其他機器。

4-1. 設定Private key 登入設定

使用 ssh-keygen 產生key pair時會詢問你一組密碼(實際上你可以偷懶使用空白密碼),然後再透過ssh-copy-id這個tool 幫你把key 傳到其他台機器

# ssh-keygen -t rsa -f ~/.ssh/id_rsa -b 4096 -C “iamcomment”
# ssh-copy-id -i .ssh/id_rsa.pub root@lab-hadoop-m2
# ssh-copy-id -i .ssh/id_rsa.pub root@lab-hadoop-m3

4-2. 設定slave

編輯/usr/lib/spark/conf/slaves 這個檔案,輸入slave 的ip或是hostname,因為我m1那台機器上也想跑一個spark worker所以我的設定如下:

localhost
lab-hadoop-m1
lab-hadoop-m2

4-3. 設定/usr/lib/spark/conf/spark-env.sh ,根據你的需求去設定裡面的內容,請參考cluster-launch-scriptsconfiguration

#!/usr/bin/env bash

# This file contains environment variables required to run Spark. Copy it as
# spark-env.sh and edit that to configure Spark for your site. At a minimum,
# the following two variables should be set:
# - SCALA_HOME, to point to your Scala installation, or SCALA_LIBRARY_PATH to
#   point to the directory for Scala library JARs (if you install Scala as a
#   Debian or RPM package, these are in a separate path, often /usr/share/java)
# - MESOS_NATIVE_LIBRARY, to point to your libmesos.so if you use Mesos
#
# If using the standalone deploy mode, you can also set variables for it:
# - SPARK_MASTER_IP, to bind the master to a different IP address
# - SPARK_MASTER_PORT / SPARK_MASTER_WEBUI_PORT, to use non-default ports
# - SPARK_WORKER_CORES, to set the number of cores to use on this machine
# - SPARK_WORKER_MEMORY, to set how much memory to use (e.g. 1000m, 2g)
# - SPARK_WORKER_PORT / SPARK_WORKER_WEBUI_PORT
# - SPARK_WORKER_INSTANCES, to set the number of worker instances/processes
#   to be spawned on every slave machine

SPARK_MASTER_WEBUI_PORT=8082
SPARK_WORKER_MEMORY=1g


5. 啟動Spark Cluster

啟動所有Master 和 Slave(Worker)

/usr/lib/spark/bin/spark-all.sh

如果沒有意外,就會看到以下訊息,Master在lab-hadoop-m1啟動,然後Slave 分別在lab-hadoop-m2,lab-hadoop-m3 啟動
starting spark.deploy.master.Master, logging to /usr/lib/spark-0.7.3/bin/../logs/spark-root-spark.deploy.master.Master-1-lab-hadoop-m1.out
Master IP: lab-hadoop-m1
cd /usr/lib/spark-0.7.3/bin/.. ; /usr/lib/spark/bin/start-slave.sh 1 spark://lab-hadoop-m1:7077
localhost: starting spark.deploy.worker.Worker, logging to /usr/lib/spark-0.7.3/bin/../logs/spark-root-spark.deploy.worker.Worker-1-lab-hadoop-m1.out
lab-hadoop-m3: starting spark.deploy.worker.Worker, logging to /usr/lib/spark-0.7.3/bin/../logs/spark-root-spark.deploy.worker.Worker-1-lab-hadoop-m3.out
lab-hadoop-m2: starting spark.deploy.worker.Worker, logging to /usr/lib/spark-0.7.3/bin/../logs/spark-root-spark.deploy.worker.Worker-1-lab-hadoop-m2.out


6. 測試Spark 是否安裝順利

在/usr/lib/spark 裡面執行以下指令

./run spark.examples.SparkLR local[2]


如果有跑出以下結果,代表安裝順利

Final w: (5816.075967498865, 5222.008066011391, 5754.751978607454, 3853.1772062206846, 5593.565827145932, 5282.387874201054, 3662.9216051953435, 4890.78210340607, 4223.371512250292, 5767.368579668863)


7. 檢視Web 管理介面

在瀏覽器輸入lab-hadoop-m1的ip,port 8020,應該就可以看到以下畫面



好Spark 裝好了,接下來就換Shark~

延伸閱讀:

[1] Shark 的實驗筆記

2013年8月8日 星期四

第一次玩Spark Shark 就上手 - 簡介篇

圖片來源:自行整理 (這隻鯊魚看起來比較威~XD)

這次要介紹的就在上一篇 Hadoop / Haddop like framework and ecosystem project 文章中有提到的Shark 和 Spark。

SharkSpark 都是屬於  Berkeley amplab BDAS. the Berkeley Data Analytics Stack  中的子系統,BDAS 的目標就是要打造一套與現有Hadoop 相容,卻又速度更快,更方便使用的系統,然後每一個子系統都可以單獨運作,整個BDAS架構如下圖所示:



BDAS是 Berkeley amplab 負責執行總經費達三千萬美金的龐大計畫(不知道對美國來說算是普通而已?),經費來源一半政府,一半業界 ,預計執行六年,目前已經執行到一半,已經產出相對穩定的專案有三個:
  

Mesos


整個BDAS的最底層用來管理Cluster 的Framework 地位等同於Hadoop2 的 YARN,特色就是可以用來管理各種不同的Cluster 包含Hadoop、Spark、MPI...等。

(目前已經是Apache 的top-level project)

Spark 


有別於Hadoop 的Disk-based MapReduce,Spark 強調的是in-memory cluster computing,號稱比Hadoop快100倍,並且有以下特點:

1. 是以Scala 撰寫 (但是有支援Scala Java 和 Python)
2. 以Resilient Distributed Dataset (RDD) 的方式達到Distributed memory layer for sharing.
3. Compatible with Hadoop Storage API,可以無縫介接Hdfs 、S3...等

(目前還在Apache incubator 階段) 

Shark


架構在Spark之上用來取代Hive ,也號稱比Hive快上30 倍以上,主要的特色就是Compatible with Hive 語法,Shark 可以直接吃Hive 的metastore 資料,也可以直接下HiveQL 去query 資料。


接下來我預計會陸續整理筆記(如果我沒偷懶...)


    PS. 有空可以去看Spark 和Shark 的Source code ...目前市面上聽到新架構新套件他全都用上了akka、spray、netty...等,絕對令人大開眼界...:P


    延伸閱讀:

    [1] An Introduction to the Berkeley Data Analytics Stack (BDAS) Featuring Spark, Spark Streaming, and Shark
    [2] Introduction to Spark, Shark, BDAS and AMPLab
    [3] Transforming Big Data with Spark and Shark - AWS Re:Invent 2012 BDT 305
    [4] Shark: Real-time queries and analytics for big data 
    [5] Spark隨談 (一) - 總體架構