從一台電腦到一個叢集:Big Data、MapReduce、Hadoop 與 Spark 到底在解決什麼問題?

發布日期:

提到 Big Data,初學者的第一個反應可能是:就是很多很多的資料吧。

但如果只是資料很多,問題真的有這麼複雜到我們需要為此發展出許多框架嗎?

假設今天只有 1,000 筆資料,我們可能只需要一台電腦、一個程式,再加上一份檔案,就可以很輕鬆地完成分析。

但如果資料從 1,000 筆變成 1 PB 呢?

這時候問題就不只是資料很多的問題而已,而是我們開始需要思考:

  • 一台電腦放得下嗎?
  • 一台電腦算得完嗎?
  • 如果資料分散在很多台機器上,要怎麼處理?
  • 如果其中一台機器突然壞掉怎麼辦?
  • 不同機器處理完之後,要怎麼把結果合併起來?

這也是 Big Data Systems 想要解決的問題。

這篇文章將從最基本的 Distributed Computing 開始,一路介紹 MapReduce、Hadoop 與 Spark,最後用一個實際的 Word Count Demo,把這些概念串起來。


Big Data 想解決什麼問題?

想像一下,你是一名 Netflix 工程師。

每天,全球的使用者都會產生大量觀看紀錄。假設某一天主管突然問你:

昨天哪部影集被觀看最多次?

如果只有幾千筆觀看紀錄,這其實是一個非常簡單的問題。

我們只需要讀取 log 檔,寫一個簡單的 script 計算每部影集出現的次數,排序後再找出觀看次數最高的那部就完成了。

一台電腦就可以完成。

但如果資料規模已經大到 PB 等級,事情就完全不一樣了。

單一機器可能會遇到:

  • Memory 不夠
  • Disk I/O 成為瓶頸
  • CPU 處理時間太長
  • 單一機器故障造成整個工作失敗

這時候,「買一台更強的電腦」或「幫電腦裝更多記憶體」未必是最好的答案(其中一個原因當然會是現在記憶體超貴)。


Scale Up 還是 Scale Out?

面對單機效能不足,我們大致有兩種思考方式。

Scale Up:買更大的電腦

也就是 Vertical Scaling。

例如原本:

4 CPU
16 GB RAM
1 TB SSD

我們升級成:

64 CPU
512 GB RAM
10 TB SSD

讓單一機器變得更強。

這種方式很直覺,但硬體終究會有上限,而且當資料規模持續成長時,成本也可能快速增加。


Scale Out:增加更多電腦

另一種方式則是 Horizontal Scaling。

例如原本只有一個 Node,現在變成有四個 Node

我們就可以把工作分散到多台配置相同機器上。

這就是 Distributed Computing 的核心概念之一:

Don’t buy a bigger computer. Use more computers.

與其一直想辦法讓一台電腦變得更強,不如讓更多台電腦一起工作。


Divide & Conquer

那麼問題來了,如果我們有一份 100 GB 的資料,要怎麼讓四台機器一起處理?

直覺的方法是

                 100 GB
                    │
          ┌─────────┼─────────┬─────────┐
          ▼         ▼         ▼         ▼
        25 GB     25 GB     25 GB     25 GB
          │         │         │         │
          ▼         ▼         ▼         ▼
        Node 1    Node 2    Node 3    Node 4

把一個大型問題拆成多個小問題,再讓不同機器平行處理。

這就是 Divide & Conquer

先 Split,再 Process,最後 Aggregate。

這個概念很簡單,但實際上真正困難是切開之後,誰負責管理?資料怎麼傳?機器壞掉怎麼辦?最後怎麼把結果合起來?


Distributed System 的挑戰

當我們從單機走向 Cluster,馬上會遇到幾個問題。

Data:資料要怎麼分?

100 GB 的資料要怎麼切?

誰負責哪一部分?


Network:機器之間怎麼傳資料?

如果 Node 1 算出來的資料需要交給 Node 3 繼續運算,那資料要怎麼傳?

如果資料量非常大,Network 也可能成為新的瓶頸。


Failure:其中一台掛掉怎麼辦?

假設四台機器正在一起處理資料

Node 3 掛掉之後,整個工作是不是都要重新開始?

如果沒有額外的機制,Node 3 原本負責的那份資料就會遺失

算出來的結果就會是錯誤的。


Coordination:最後怎麼把結果合起來?

每台機器都算完自己的資料之後,誰負責把結果整合起來?


這些問題就是 Distributed Computing 真正複雜的地方。

因此,我們需要 Framework 幫我們處理這些複雜性。

而 Hadoop、Spark 等工具,就是在這個背景下出現的。


MapReduce:把大型工作拆成兩個步驟

MapReduce 是一種非常經典的分散式資料處理模型。

它把工作拆成 Map → Shuffle → Reduce。

Map 負責產生 Key-Value,Reduce 負責把相同 Key 的資料聚合起來。


Word Count 是非常經典的 MapReduce 範例。

假設我們有一個句子:Big Data is Big.

我們想知道每個單字出現幾次。

Map

Map 階段可以把資料轉換成:

Big   → 1
Data  → 1
is    → 1
Big   → 1

這時候還沒有真正完成統計。

接下來就進入 Shuffle。


Shuffle

Shuffle 會把相同 Key 的資料放在一起:

Big   → [1, 1]
Data  → [1]
is    → [1]

最後交給 Reduce。


Reduce

Reduce 只需要把相同 Key 的數字加起來:

Big   → 3
Data  → 2
is    → 1

Word Count 就完成了。

整個流程可以簡化成

Input Map Shuffle Reduce Output


MapReduce 最大的優勢:Parallelism

MapReduce 最重要的概念之一,就是可以平行處理。

假設我們有四份資料:

File
 │
 ├── Part 1 → Mapper 1
 ├── Part 2 → Mapper 2
 ├── Part 3 → Mapper 3
 └── Part 4 → Mapper 4

四個 Mapper 可以同時處理不同資料。

這就是我們前面提到的 Divide & Conquer,把大型工作拆成很多小工作,再平行處理。

如果資料規模持續增加,我們也可以增加更多 Worker。

這就是 Distributed Computing 的價值。


Hadoop:讓 MapReduce 跑在 Cluster 上

MapReduce 的概念很好,但我們還需要真正把它放到很多台機器上。

這就是 Hadoop 發揮作用的地方。

Hadoop 不只是 MapReduce。

其中一個重要組成部分是:

HDFS:Hadoop Distributed File System

HDFS 負責的是資料的存儲

而 MapReduce 則負責資料的處理與運算


HDFS:大型檔案怎麼存?

假設我們有一個大型檔案。

HDFS 不會把它當成一個完整的檔案,只放在某一台機器上。

它會把大型檔案切成多個 Block(預設是 128 MB)。

例如:

100 GB File
     │
     ▼
 ┌──────┬──────┬──────┬──────┐
 │25 GB │25 GB │25 GB │25 GB │
 └──────┴──────┴──────┴──────┘
     │      │      │      │
     ▼      ▼      ▼      ▼
   Node1  Node2  Node3  Node4

這樣我們就可以把資料分散到不同 DataNode。


如果資料只被分散儲存,還有一個問題:如果某一台機器壞掉呢?

所以 HDFS 還有另一個重要概念:Replication

也就是讓一個 Block 保留多份副本(類似 Kubernetes 的ReplicaSet)。

例如:

Block A
   │
   ├── DataNode 1
   ├── DataNode 2
   └── DataNode 3

當其中一台 DataNode 發生故障時,其他副本仍然可以使用。


Hadoop 中幾個專有名詞

NameNode

NameNode 負責管理 HDFS 的 Metadata。

NameNode 不負責實際儲存這些資料。


DataNode

DataNode 才是真正儲存和使用 Block 的 Worker。


Hadoop WordCount

我準備了一個可以直接執行的 Hadoop 範例

完整程式碼放在 GitHub:GitHub Repository

其中 Hadoop Demo 放在 hadoop/,使用 Docker 建立 Hadoop Cluster,包含 HDFS、YARN,以及 Hadoop Streaming WordCount。

執行:

ShellScript
bash hadoop/run_docker.sh

也可以指定 Worker 數量:

ShellScript
bash hadoop/run_docker.sh hadoop/wordcount/sample_input.txt 3

啟動之後,可以透過 NameNode 與 YARN 的 Web UI 觀察 Cluster 狀態,也可以直接從 terminal 輸出看到 word count 結果。

這個 Demo 使用 Hadoop Streaming。

Mapper 負責把每個單字轉成:

Markdown
word    1

Reducer 則負責把相同的 word 聚合起來。

Markdown
word    [1, 1, 1] => word    3 

這其實就是我們前面介紹的 MapReduce。

前面我們是在理解演算法,現在 Hadoop 則是真的把它放到 Cluster 裡執行。


到這裡看起來 Hadoop 已經非常厲害了。

資料可以分散,工作可以平行,機器壞掉也可以透過 Replication 等機制處理。

但如果我們要連續做很多次運算呢?

傳統 MapReduce 的處理方式會涉及大量 Disk I/O。

資料被讀取、處理、寫回,再被下一個 Job 讀取。

當我們需要進行大量、連續的資料運算時,這可能會成為效能瓶頸。


Spark:把資料處理變得更快、更容易

這也是 Apache Spark 出現的重要背景之一。

Spark 的核心特色之一是 In-memory computing

簡單來說,就是盡可能把需要重複使用的資料保留在 Memory 中,而不是每次都重新從 Disk 讀取。

因此,如果工作流程需要對同一批資料進行多次運算,Spark 就能避免一些不必要的 Disk I/O。


Spark WordCount

Spark 也可以用非常類似的方式完成 Word Count。

Input Read Data RDD Partitions Map (word, 1) ReduceByKey Word Count

在 Spark Demo 中,我們使用 PySpark 完成這個流程。

程式會:

  1. 將文字讀入 Spark RDD
  2. 將每個單字轉成 (word, 1)
  3. 使用 reduceByKey 聚合
  4. 排序並輸出結果

完整 Demo 同樣放在 GitHub repository 中。


那 Hadoop 和 Spark 差在哪?

Spark 並不是 Hadoop 的「下一個版本」。

兩者解決的問題有部分重疊,但角色並不完全相同。

可以先用一個非常簡單的方式理解:

HDFSMapReduceSpark
主要用途StorageProcessingProcessing
核心概念Distributed StorageMap / Shuffle / ReduceIn-memory Processing
資料處理BatchBatch + 更多處理模型
典型特色分散式儲存平行批次運算更適合需要重複運算的工作

HDFS 解決「資料放在哪裡」的問題。

MapReduce 解決「資料怎麼分散式處理」的問題。

Spark 則提供另一套更適合許多資料處理工作的分散式計算引擎。

這也是為什麼在 Big Data Ecosystem 裡,我們不應該把 Hadoop、MapReduce、HDFS、Spark 當成完全相同的東西。


回到一開始的 Netflix 問題

現在回到一開始的問題:「昨天哪部影集被觀看最多次?」

如果資料只有幾千筆,一台機器運算完全沒問題。

但如果資料已經到了 PB 等級,我們可以開始考慮使用框架來處理,像是

                 Big Data
                    │
                    ▼
              Distributed
               Computing
                    │
          ┌─────────┴─────────┐
          ▼                   ▼
       Storage             Processing
          │                   │
          ▼                   ▼
        HDFS               MapReduce
                              │
                              ▼
                            Spark

資料可以被拆分、分散儲存,再交給多台機器平行處理。

原本一台電腦很難處理的工作,就可以被轉換成一個 Distributed System 的問題。


Big Data Ecosystem

所以當我們講 Big Data 時,真正重要的是理解它們各自解決什麼問題。

再往真實世界延伸,Big Data Ecosystem 還會包含更多不同工具,分別處理資料收集、訊息傳遞、儲存、計算、分析與視覺化等問題。

但這篇文章最重要的,是先建立一個最基本的 mental model:

資料太大 → 一台機器不夠 → 分散到多台機器 → 需要 Framework 幫我們管理這個 Distributed System。


Takeaway

Scale Out

Big Data 的核心思想之一,是 Big Data → More Machines

當單一機器成為瓶頸時,不一定只能一直買更大的機器。

我們也可以透過增加更多機器來擴展系統。


Divide & Conquer

大型問題可以拆成:Split Process Aggregate
今天介紹的 MapReduce 就是非常經典的例子。


Frameworks Matter

真正困難的不是「把資料切成幾份」。

而是:

  • 資料怎麼分?
  • 機器之間怎麼傳?
  • 機器掛掉怎麼辦?
  • 結果怎麼合併?
  • 如何讓整個 Cluster 協同工作?

Hadoop、Spark 等 Framework 的價值,就是幫我們處理這些 Distributed Computing 的複雜性。


結語

大數據背後的問題非常單純,那就是

當一台電腦已經無法處理我們的資料時,我們該怎麼辦?

答案就是從一台機器走向多台機器。

接著再透過 Distributed Computing,把原本一個大型問題拆成許多小問題,讓不同機器平行處理

而 Hadoop、MapReduce 與 Spark,則是在這個過程中,分別提供不同的儲存與計算能力。

當資料大到單機無法有效處理時,我們如何讓整個系統一起工作。


Demo Repository

如果想實際操作文章中的範例,可以參考我的 GitHub Repository:

Big Data Demos — GitHub

Repository 包含三個部分:

  • MapReduce:用 Python 從最基本的 Word Count 開始理解 Map、Shuffle、Reduce
  • Hadoop:使用 Docker 建立 Hadoop Cluster,實際執行 Hadoop Streaming Word Count
  • Spark:使用 PySpark 執行 Spark Word Count

其中 MapReduce Demo 也提供大型輸入資料產生工具,可以用來觀察單機處理大量資料時可能遇到的問題;Hadoop 與 Spark Demo 則可以直接啟動 Cluster 進行實際操作。

留言

發表留言