1. MapReduce入門
1.1 MapReduce的思想
MapReduce的思想核心是"分而治之" , 適用于大量的復雜的任務處理場景 (大規模資料處理場景) .
Map負責"分" , 即把復雜的任務分解為若干個"簡單的任務"來進行處理. 可以進行拆分的前提是這些小任務并行計算, 彼此間幾乎沒有依賴關系.
Reduce負責"合" , 即對map階段的結果進行全域匯總.
這兩個階段合起來正是MR思想的體現.
1.2 MapReduce設計構思
MapReduce是一個分部式運算程式的編程框架, 核心功能是將用戶撰寫的業務邏輯代碼和自帶默認組件整合成一個完整的分布式運算程式. 并發運行在Hadoop集群上.
既然是做計算的框架, 那么表現形式就是有個輸入 (input) , MR操作這個輸入, 通過本身定義好的計算模型, 得到一個輸出 (output) .
MR就是一種簡化并行計算的編程模型, 降低了開發并行應用的入門門檻.
Hadoop MapReduce構思體現在三個方面:
-
- 如何對付大資料處理: 分而治之
對相互間不具有計算依賴關系的大資料, 實作并行最自然的方法就是采取分而治之的策略. 并行計算的第一個重要問題是如何劃分計算任務或者計算資料以便對劃分的子任務或資料塊同時進行計算. 不可分拆的計算任務或相互間有依賴關系的資料無法進行并行計算.
-
- 構建抽象模型: Map和Reduce
MR借鑒了函式式語言中的思想, 用Map和Reduce兩個函式提供了高層的并行編程抽象模型.
Map: 對一組資料元素進行某種重復式的處理.
Reduce: 對Map的中間結果進行某種進一步的結果整理.
MapReduce中定義了如下的Map和Reduce兩個抽象的編程介面, 由用戶去編程實作:
Map: (k1, v1) -> [(k2, v2)]
Reduce: (k2, [v2]) -> [(k3, v3)]
Map和Reduce為我們提供了一個清晰的操作介面抽象描述. 通過以上兩個編程介面, 可以看出MR處理的資料型別是<key, value>鍵值對.
-
- 統一架構, 隱藏系統層細節
如何提供統一的計算框架, 如果沒有統一封裝底層細節, 那么我們則需要考慮諸如資料存盤, 劃分, 分發, 結果手機, 錯誤恢復等諸多細節. 為此, MR設計并提供了統一的計算框架, 隱藏了絕大多數系統層面的處理細節.
MapReduce最大的亮點在于通過抽象模型和計算框架把需要做什么 (what need to do) 與具體怎么做 (how to do) 分開了, 為我們提供一個抽象和高層的編程介面和框架. 我們僅需要關心其應用層的具體計算問題, 僅需撰寫少量的處理應用本身計算問題的程式代碼. 如何具體完成這個并行計算任務所相關的諸多系統層細節被隱藏起來, 交給計算框架去處理: 從分布代碼的執行, 到大到數千小到單個節點集群的自動調度使用.
1.3 MapReduce框架結構
一個完整的MR程式在分布式運行時有三類實體行程:
1) MRAppMaster: 負責整個程式的程序調度及狀態協調.
2) MapTask: 負責Map階段的整個資料處理流程.
3) ReduceTask: 負責Reduce階段的整個資料處理流程.


1.4 Map, Reduce, Split總結

2. MapReduce編程規范及示例
2.1 編程規范
開發步驟一共8步
-
- MapTask階段2步
1) 設定InputFormat (通常使用TextInputFormat) 的型別和資料的路徑 -- 獲取資料的程序 (可以得到K1, V1) .
2) 自定義Mapper -- 將K1, V1轉為K2, V2.
-
- shuffle階段4步
3) 磁區的動作, 如果有多個Reduce才去考慮磁區, 默認只有一個Reduce, 磁區可以省略.
4) 排序, 默認對K2進行排序 (字典序) -- 管好K2就行.
5) 規約, combiner是一個區域的Reduce, Map端的合并, 是對MR的優化操作, 不會影響任何結果, 減少網路傳輸, 默認可以省略.
6) 分組, 相同的K (K2) 對應的V會放到同一個集合中 -- 將Map傳遞的K2, V2變成新的K2, V2.
-
- Reduce階段2步
7) 自定義Reducer得到K2, V2轉為K3, V3.
8) 設定OutputFormat和資料的路徑 -- 生成結果檔案.

2.2 WordCount案例

3. MapReduce程式運行模式
3.1 本地運行模式
1) MR程式是被提交給LocalJobRunner在本地以單行程的形式運行.
2) 而處理的資料及輸出結果可以在本地檔案系統, 也可以在HDFS上.
3) 怎么樣實作本地運行? 寫一個程式, 不要帶集群的組態檔
本質是程式的conf中是否有mapreduce.framework.name=local以及yarn.resourcemanager.hostname引數
4) 本地模式非常便于進行業務邏輯的debug, 只要打斷點即可.
3.2 集群運行模式
1) 將MR程式提交給yarn集群, 分發到很多的節點上并發執行.
2) 處理的資料和輸出結果應該位于HDFS檔案系統.
3) 提交集群的實作步驟:
將程式打成jar包, 然后在集群的任意一個節點上用Hadoop命令啟動:
hadoop jar wordcount.jar cn.itcast.mr.wordcount.WordCountRunner args
4. 深入MapReduce
4.1 MapReduce的輸入和輸出
MR框架運轉在<key, value>鍵值對上, 也就是說, 框架把作業的輸入看成是一組<key, value>鍵值對, 同樣也產生一組<key, value>鍵值對作為作業的輸出, 這兩組鍵值對可能是不同的.
一個MR作業的輸入和輸出型別如下圖所示: 可以看出在整個標準流程中, 會有三組<key, value>鍵值對型別的存在.

4.2 Mapper任務執行程序詳解
第一階段: 把輸入目錄下檔案按照一定的標準逐個進行邏輯切片, 形成切片規劃. 默認情況下, Split size = Block size. 每一個切片由一個MapTask處理. (getSplits)
第二階段: 對切片中的資料按照一定的規則決議成<key, value>對. 默認規則是把每一行文本內容決議成鍵值對. key是每一行的起始位置 (單位是位元組) , value是本行的文本內容. (TextInputFormat)
第三階段: 呼叫Mapper類中的map方法, 上階段中每決議出來的一個<k, v>, 呼叫一次map方法. 每次呼叫map方法會輸出零個或多個鍵值對.
第四階段: 按照一定的規則對第三階段輸出的鍵值對進行磁區. 默認是只有一個區. 磁區的數量就是Reducer任務運行的數量. 默認只有一個Reducer任務.
第五階段: 對每個磁區中的鍵值對進行排序. 首先, 按照鍵進行排序, 對于鍵相同的鍵值對, 按照值進行排序. 比如三個鍵值對<2, 2>, <1, 3>, <2, 1>, 鍵和值分別是整數. 那么排序后的結果是<1, 3>, <2, 1>, <2, 2>. 如果有第六階段, 那么進入第六階段, 如果沒有, 直接輸出到檔案中.
第六階段: 對資料進行區域聚合處理, 也就是combiner處理. 鍵相等的鍵值對會呼叫一次reduce方法. 經過這一階段, 資料量會減少. 本階段默認是沒有的.
4.3 Reducer任務執行程序詳解
第一階段: Reducer任務會主動從Mapper任務復制其輸出的鍵值對. Mapper任務可能會有很多, 因此Reducer會復制多個Mapper的輸出.
第二階段: 把復制到Reducer本地資料, 全部進行合并, 即把分散的資料合并成一個大的資料. 再對合并后的資料排序.
第三階段: 對排序后的鍵值對呼叫reduce方法. 鍵相等的鍵值對呼叫一次reduce方法, 每次呼叫會產生零個或者多個鍵值對. 最后把這些輸出的鍵值對寫入到HDFS檔案中.
在整個MR程式的開發程序中, 我們最大的作業是覆寫map函式和覆寫reduce函式.
5. MapReduce的序列化
5.1 概述
序列化是指把結構化物件轉化為位元組流.
反序列化是序列化的逆程序. 把位元組流轉化為結構化物件.
當要在行程間傳遞物件或持久化物件的時候, 就需要序列化物件成位元組流, 反之當要將接收到或從磁盤讀取的位元組流轉換成物件, 就要進行反序列化.
Java序列化是一個重量級序列化框架, 一個物件被序列化后, 會附帶很多額外的資訊, 不便于在網路中高效傳輸. 所以, Hadoop自己開發了一套序列化機制 (Writable) , 不用像Java物件類一樣傳輸多層的父子關系, 需要哪個屬性就傳輸哪個屬性值, 大大的減少網路傳輸的開銷.
Writable是Hadoop的序列化格式, Hadoop定義了這樣一個Writable介面. 一個類要支持可序列化只需要實作這個介面即可.
public interface Wriable {
void wirte (DataOutput out) throws IOException;
void readFields (DataInput in) throws IOException;
}
5.2 Writable序列化介面
如需將自定義的bean放在key中傳輸, 則還需要實作comparable介面, 因為MR框中的shuffle程序一定會對key進行排序, 此時, 自定義的bean實作的介面應該是:
public class FlowBean implements WritableComparable<FlowBean>
compareTo方法用于將當前物件與方法的引數進行比較:
-
- 如果指定的數與引數相等回傳 0
- 如果指定的數小于引數回傳 -1
- 如果指定的數大于引數回傳 1
例如: o1.compareTo(o2)
回傳正數的話, 當前物件 (呼叫compareTo方法的物件o1) 要排在比較物件 (compareTo傳參物件o2) 后面, 回傳負數的話, 放在前面.

6. MapReduce的排序初步
6.1 需求
在得出統計每一個用戶 (手機號) 所耗費的總上行流量, 下行流量, 總流量結果的基礎之上再加一個需求: 將統計結果按照總流量倒序排序.
6.2 分析
基本思路: 實作自定義的bean來封裝流量資訊, 并將bean作為map輸出的key來傳輸.
MR程式在處理資料的程序中會對資料排序 (map輸出的kv對傳輸到reduce之前, 會排序), 排序的依據是map輸出的key. 所以, 我們如果要實作自己需要的排序規則, 則可以考慮將排序因素放到key中, 讓key實作介面: WritableComparable, 然后重寫key的compareTo方法.
7. MapReduce的磁區Partitioner
7.1 需求
將流量匯總統計結果按照手機歸屬地不同省份輸出到不同檔案中.
7.2 分析
Mapreduce中會將map輸出的kv對, 按照相同key分組, 然后分發給不同的reducetask.
默認的分發規則為: 根據key的hashcode%reducetask數來分發.
所以: 如果要按照我們自己的需求進行分組, 則需要改寫資料分發 (分組) 組件Partitioner, 自定義一個CustomPartitioner繼承抽象類: Partitioner, 然后在job物件中, 設定自定義partitioner: job.setPartitionerClass(CustomPartitioner.class) .
8. MapReduce的Combiner
每一個map都可能會產生大量的本地輸出, Combiner的作用就是對map端的輸出先做一次合并, 以減少在map和reduce節點之間的資料傳輸量, 以提高網路IO性能, 是MR的一種優化手段之一.
-
- Combiner是MR程式中Mapper和Reducer之外的一種組件.
- Combiner組件的父類就是Reducer.
- Combiner和reducer的區別在于運行的位置:
Combiner是在每一個maptask所在的節點運行.
Reducer是接收全域所有Mapper的輸出結果.
-
- Combiner的意義就是對每一個maptask的輸出進行區域匯總, 以減小網路傳輸量.
- 具體實作步驟:
1) 自定義一個combiner繼承Reducer, 重寫reduce方法
2) 在job中設定: job.setCombinerClass(CustomCombiner.class) .
-
- Combiner能夠應用的前提是不能影響最終的業務邏輯, 而且, Combiner的輸出kv應該跟reducer的輸入kv型別要對應起來.
轉載請註明出處,本文鏈接:https://www.uj5u.com/shujuku/39994.html
標籤:大數據
上一篇:求做過中控k28 ,h10考勤機二次開發的大神,只要能讀出用戶、考勤資料
下一篇:sybase是不是已經放棄pb了
