
文件內(nèi)容如下hello world hello spark hello mapreduce需求統(tǒng)計每個單詞出現(xiàn)多少次1、Map 階段Map 讀取每一行文本切割成單詞輸出(單詞, 1)出現(xiàn)一次單詞記計數(shù) 1輸入第一行hello world→ 切割輸出(hello,1) (world,1)第二行hello spark(hello,1) (spark,1)第三行hello mapreduce(hello,1) (mapreduce,1)Map 輸出全部 KV(hello,1) (world,1) (hello,1) (spark,1) (hello,1) (mapreduce,1)Map 做的事拆分數(shù)據(jù)打上標記多個 Map 任務并行處理不同文本塊2、Shuffle 洗牌搬運分組把相同 key 的全部數(shù)據(jù)收集到一塊發(fā)給同一個 Reduce經(jīng)過 shuffle 分揀之后分組hello → [1, 1, 1] world → [1] spark → [1] mapreduce → [1]?關(guān)鍵點所有的hello全部搬運到同一個 Reduce 任務shuffle 要磁盤讀寫 網(wǎng)絡傳輸數(shù)據(jù)傾斜就出在這里如果某個單詞幾千萬條這個 Reduce 就扛不住 OOM。3、Reduce 階段Reduce 拿到同一個 key 對應的一堆數(shù)字把數(shù)字累加求和hello: 111 3 world:1 spark:1 mapreduce:1輸出最終統(tǒng)計結(jié)果??焖儆洃汳ap拆單詞輸出 (單詞1)Shuffle把相同單詞全部匯集到一處Reduce對相同單詞的 1 累加得到總次數(shù)MapReduceMap 輸出全部寫磁盤再 shuffleSpark前面 map 操作放內(nèi)存到 shuffle 這一步才寫磁盤WordCount 單詞統(tǒng)計Map 讀取文本切分單詞輸出 key?value(word,1)Shuffle 將相同 word 的數(shù)據(jù)分組網(wǎng)絡傳輸交給同一個 ReduceReduce 對 value 集合求和輸出每個單詞計數(shù)Shuffle 是 IO 開銷最大的階段數(shù)據(jù)傾斜發(fā)生在此處