主頁 > 資料庫 > Hadoop_MapReduce_03

Hadoop_MapReduce_03

2020-09-15 00:46:27 資料庫

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了

標籤雲
其他(157675) Python(38076) JavaScript(25376) Java(17977) C(15215) 區塊鏈(8255) C#(7972) AI(7469) 爪哇(7425) MySQL(7132) html(6777) 基礎類(6313) sql(6102) 熊猫(6058) PHP(5869) 数组(5741) R(5409) Linux(5327) 反应(5209) 腳本語言(PerlPython)(5129) 非技術區(4971) Android(4554) 数据框(4311) css(4259) 节点.js(4032) C語言(3288) json(3245) 列表(3129) 扑(3119) C++語言(3117) 安卓(2998) 打字稿(2995) VBA(2789) Java相關(2746) 疑難問題(2699) 细绳(2522) 單片機工控(2479) iOS(2429) ASP.NET(2402) MongoDB(2323) 麻木的(2285) 正则表达式(2254) 字典(2211) 循环(2198) 迅速(2185) 擅长(2169) 镖(2155) 功能(1967) .NET技术(1958) Web開發(1951) python-3.x(1918) HtmlCss(1915) 弹簧靴(1913) C++(1909) xml(1889) PostgreSQL(1872) .NETCore(1853) 谷歌表格(1846) Unity3D(1843) for循环(1842)

熱門瀏覽
  • GPU虛擬機創建時間深度優化

    **?桔妹導讀:**GPU虛擬機實體創建速度慢是公有云面臨的普遍問題,由于通常情況下創建虛擬機屬于低頻操作而未引起業界的重視,實際生產中還是存在對GPU實體創建時間有苛刻要求的業務場景。本文將介紹滴滴云在解決該問題時的思路、方法、并展示最終的優化成果。 從公有云服務商那里購買過虛擬主機的資深用戶,一 ......

    uj5u.com 2020-09-10 06:09:13 more
  • 可編程網卡芯片在滴滴云網路的應用實踐

    **?桔妹導讀:**隨著云規模不斷擴大以及業務層面對延遲、帶寬的要求越來越高,采用DPDK 加速網路報文處理的方式在橫向縱向擴展都出現了局限性。可編程芯片成為業界熱點。本文主要講述了可編程網卡芯片在滴滴云網路中的應用實踐,遇到的問題、帶來的收益以及開源社區貢獻。 #1. 資料中心面臨的問題 隨著滴滴 ......

    uj5u.com 2020-09-10 06:10:21 more
  • 滴滴資料通道服務演進之路

    **?桔妹導讀:**滴滴資料通道引擎承載著全公司的資料同步,為下游實時和離線場景提供了必不可少的源資料。隨著任務量的不斷增加,資料通道的整體架構也隨之發生改變。本文介紹了滴滴資料通道的發展歷程,遇到的問題以及今后的規劃。 #1. 背景 資料,對于任何一家互聯網公司來說都是非常重要的資產,公司的大資料 ......

    uj5u.com 2020-09-10 06:11:05 more
  • 滴滴AI Labs斬獲國際機器翻譯大賽中譯英方向世界第三

    **桔妹導讀:**深耕人工智能領域,致力于探索AI讓出行更美好的滴滴AI Labs再次斬獲國際大獎,這次獲獎的專案是什么呢?一起來看看詳細報道吧! 近日,由國際計算語言學協會ACL(The Association for Computational Linguistics)舉辦的世界最具影響力的機器 ......

    uj5u.com 2020-09-10 06:11:29 more
  • MPP (Massively Parallel Processing)大規模并行處理

    1、什么是mpp? MPP (Massively Parallel Processing),即大規模并行處理,在資料庫非共享集群中,每個節點都有獨立的磁盤存盤系統和記憶體系統,業務資料根據資料庫模型和應用特點劃分到各個節點上,每臺資料節點通過專用網路或者商業通用網路互相連接,彼此協同計算,作為整體提供 ......

    uj5u.com 2020-09-10 06:11:41 more
  • 滴滴資料倉庫指標體系建設實踐

    **桔妹導讀:**指標體系是什么?如何使用OSM模型和AARRR模型搭建指標體系?如何統一流程、規范化、工具化管理指標體系?本文會對建設的方法論結合滴滴資料指標體系建設實踐進行解答分析。 #1. 什么是指標體系 ##1.1 指標體系定義 指標體系是將零散單點的具有相互聯系的指標,系統化的組織起來,通 ......

    uj5u.com 2020-09-10 06:12:52 more
  • 單表千萬行資料庫 LIKE 搜索優化手記

    我們經常在資料庫中使用 LIKE 運算子來完成對資料的模糊搜索,LIKE 運算子用于在 WHERE 子句中搜索列中的指定模式。 如果需要查找客戶表中所有姓氏是“張”的資料,可以使用下面的 SQL 陳述句: SELECT * FROM Customer WHERE Name LIKE '張%' 如果需要 ......

    uj5u.com 2020-09-10 06:13:25 more
  • 滴滴Ceph分布式存盤系統優化之鎖優化

    **桔妹導讀:**Ceph是國際知名的開源分布式存盤系統,在工業界和學術界都有著重要的影響。Ceph的架構和演算法設計發表在國際系統領域頂級會議OSDI、SOSP、SC等上。Ceph社區得到Red Hat、SUSE、Intel等大公司的大力支持。Ceph是國際云計算領域應用最廣泛的開源分布式存盤系統, ......

    uj5u.com 2020-09-10 06:14:51 more
  • es~通過ElasticsearchTemplate進行聚合~嵌套聚合

    之前寫過《es~通過ElasticsearchTemplate進行聚合操作》的文章,這一次主要寫一個嵌套的聚合,例如先對sex集合,再對desc聚合,最后再對age求和,共三層嵌套。 Aggregations的部分特性類似于SQL語言中的group by,avg,sum等函式,Aggregation ......

    uj5u.com 2020-09-10 06:14:59 more
  • 爬蟲日志監控 -- Elastc Stack(ELK)部署

    傻瓜式部署,只需替換IP與用戶 導讀: 現ELK四大組件分別為:Elasticsearch(核心)、logstash(處理)、filebeat(采集)、kibana(可視化) 下載均在https://www.elastic.co/cn/downloads/下tar包,各組件版本最好一致,配合fdm會 ......

    uj5u.com 2020-09-10 06:15:05 more
最新发布
  • day02-2-商鋪查詢快取

    功能02-商鋪查詢快取 3.商鋪詳情快取查詢 3.1什么是快取? 快取就是資料交換的緩沖區(稱作Cache),是存盤資料的臨時地方,一般讀寫性能較高。 快取的作用: 降低后端負載 提高讀寫效率,降低回應時間 快取的成本: 資料一致性成本 代碼維護成本 運維成本 3.2需求說明 如下,當我們點擊商店詳 ......

    uj5u.com 2023-04-20 08:33:24 more
  • MySQL中binlog備份腳本分享

    關于MySQL的二進制日志(binlog),我們都知道二進制日志(binlog)非常重要,尤其當你需要point to point災難恢復的時侯,所以我們要對其進行備份。關于二進制日志(binlog)的備份,可以基于flush logs方式先切換binlog,然后拷貝&壓縮到到遠程服務器或本地服務器 ......

    uj5u.com 2023-04-20 08:28:06 more
  • day02-短信登錄

    功能實作02 2.功能01-短信登錄 2.1基于Session實作登錄 2.1.1思路分析 2.1.2代碼實作 2.1.2.1發送短信驗證碼 發送短信驗證碼: 發送驗證碼的介面為:http://127.0.0.1:8080/api/user/code?phone=xxxxx<手機號> 請求方式:PO ......

    uj5u.com 2023-04-20 08:27:27 more
  • 快取與資料庫雙寫一致性幾種策略分析

    本文將對幾種快取與資料庫保證資料一致性的使用方式進行分析。為保證高并發性能,以下分析場景不考慮執行的原子性及加鎖等強一致性要求的場景,僅追求最終一致性。 ......

    uj5u.com 2023-04-20 08:26:48 more
  • sql陳述句優化

    問題查找及措施 問題查找 需要找到具體的代碼,對其進行一對一優化,而非一直把關注點放在服務器和sql平臺 降低簡化每個事務中處理的問題,盡量不要讓一個事務拖太長的時間 例如檔案上傳時,應將檔案上傳這一步放在事務外面 微軟建議 4.啟動sql定時執行計劃 怎么啟動sqlserver代理服務-百度經驗 ......

    uj5u.com 2023-04-20 08:26:35 more
  • 云時代,MySQL到ClickHouse資料同步產品對比推薦

    ClickHouse 在執行分析查詢時的速度優勢很好的彌補了MySQL的不足,但是對于很多開發者和DBA來說,如何將MySQL穩定、高效、簡單的同步到 ClickHouse 卻很困難。本文對比了 NineData、MaterializeMySQL(ClickHouse自帶)、Bifrost 三款產品... ......

    uj5u.com 2023-04-20 08:26:29 more
  • sql陳述句優化

    問題查找及措施 問題查找 需要找到具體的代碼,對其進行一對一優化,而非一直把關注點放在服務器和sql平臺 降低簡化每個事務中處理的問題,盡量不要讓一個事務拖太長的時間 例如檔案上傳時,應將檔案上傳這一步放在事務外面 微軟建議 4.啟動sql定時執行計劃 怎么啟動sqlserver代理服務-百度經驗 ......

    uj5u.com 2023-04-20 08:25:13 more
  • Redis 報”OutOfDirectMemoryError“(堆外記憶體溢位)

    Redis 報錯“OutOfDirectMemoryError(堆外記憶體溢位) ”問題如下: 一、報錯資訊: 使用 Redis 的業務介面 ,產生 OutOfDirectMemoryError(堆外記憶體溢位),如圖: 格式化后的報錯資訊: { "timestamp": "2023-04-17 22: ......

    uj5u.com 2023-04-20 08:24:54 more
  • day02-2-商鋪查詢快取

    功能02-商鋪查詢快取 3.商鋪詳情快取查詢 3.1什么是快取? 快取就是資料交換的緩沖區(稱作Cache),是存盤資料的臨時地方,一般讀寫性能較高。 快取的作用: 降低后端負載 提高讀寫效率,降低回應時間 快取的成本: 資料一致性成本 代碼維護成本 運維成本 3.2需求說明 如下,當我們點擊商店詳 ......

    uj5u.com 2023-04-20 08:24:03 more
  • day02-短信登錄

    功能實作02 2.功能01-短信登錄 2.1基于Session實作登錄 2.1.1思路分析 2.1.2代碼實作 2.1.2.1發送短信驗證碼 發送短信驗證碼: 發送驗證碼的介面為:http://127.0.0.1:8080/api/user/code?phone=xxxxx<手機號> 請求方式:PO ......

    uj5u.com 2023-04-20 08:23:11 more