跳到主要內容
Lab Grimoire
TW
請喝咖啡
動手實作

資料處理技術:批次、串流、Hadoop、Spark 與邊緣運算

批次與串流處理的差別,HDFS、MapReduce、Spark 惰性求值、Kafka、Flink 視窗與水位線、邊緣運算,以及分散式計算統計量的方法。

作者
CY
發表
AI 的資料基礎:從統計到生成式 AI · 第 7/19 篇

本系列第 7 篇|這篇回答:資料大到一台電腦處理不了時,工作要怎麼拆? 讀完這篇,你應該能:

  • 分辨批次與串流處理並配對情境
  • 說明 Hadoop、Spark、Kafka、Flink 各自解決的問題
  • 理解分散式運算為什麼要「先局部彙總再合併」

一、原理:資料大到一台機器處理不了

明細超出單機時,運算送到資料所在處,網上只傳可相加、可合併的摘要。

詳細說明

明細多到單機的記憶體與磁碟不夠用時,讀取、排序與彙總都會卡在同一台機器上。分散式處理把資料切成區塊,放到多台機器,再把運算送到資料所在的那一台,而不是把全部明細先搬去同一個地方。這個原則稱為資料局部性(data locality)。各台機器處理自己的區塊,之後只把局部結果合併。網路上走的是可相加、可合併的摘要,原始明細仍留在原地。

要等多久,是另一個選擇。批次處理(batch)先把一段時間的資料累積起來,再一次處理。它的延遲高、吞吐大。每日報表屬於這一類:不追求資料一到就更新,要的是這一段被完整算過。串流處理(stream)在資料到達時就處理,延遲低到秒級以下。即時詐欺偵測屬於這一類:判斷要落在交易進行的當下,不能等隔日那一批。可以等到隔日的完整結果用批次;必須在當下做完的判斷用串流。

資料在哪裡產生,也會改寫前半段要放在哪。邊緣運算(edge computing)在產生端附近先處理,例如工廠閘道器、攝影機或車輛。現場先做判斷或做成摘要,延遲與上傳頻寬一起下降,原始資料也不必離開現場。需要全集的工作仍放在雲端,例如集中訓練,以及把各據點的結果彙整起來。

二、方法:工具與定位

看一個工具先問它解決什麼:存放、計算、傳送,還是晚到事件如何入帳。

詳細說明

上面三個選擇各自落到不同工具。看一個工具時,先問它解決的是存放、計算、傳送,還是晚到的事件如何入帳。

Hadoop處理的是「單機放不下、也算不完」。HDFS(Hadoop Distributed File System,分散式檔案系統)把檔案切成區塊,散放在多台機器;預設每個區塊複製 3 份,用來容錯,其中一份讀不到時還有副本。MapReduce把計算分成 Map、Shuffle、Reduce:Map 在各區塊上做局部處理,Shuffle 依鍵把中間結果重新分組,Reduce 合併同一鍵的值。中間結果寫入磁碟。只跑一輪時,這個寫入還扛得住;同一份資料要反覆掃很多輪時,每一輪都讀寫磁碟,時間就堆上去。YARN負責把叢集上的 CPU 與記憶體分給這些工作。

Apache Spark把中間結果留在記憶體,迭代型工作因此不必每輪落盤。機器學習在大資料上多輪訓練,就是這種工作:Spark 通常比 MapReduce 快,差在中間結果放記憶體還是磁碟。介面上,DataFrame 與 Spark SQL 用表格描述轉換,MLlib 放機器學習演算法。Structured Streaming走微批次(micro-batch):把很短區間內到達的資料收成一小批再算。它是串流介面,執行時仍是一小批一小批。

Spark 的惰性求值(lazy evaluation)決定這些步驟何時真的跑。轉換(transformation),例如 filter、select、groupBy,只把步驟加進執行計畫(DAG),不去讀資料。遇到動作(action),例如 count、collect、write,引擎才把整條計畫放在一起最佳化並執行。因此一段只有篩選與分組、沒有任何輸出動作的程式,會很快結束,計算還沒開始。

collect() 本身是動作,而且會把全部結果拉回 driver,也就是負責協調的那一台。資料一大,driver 的記憶體不足,工作就失敗。大資料的彙總應在叢集上做完,拉回單機的只應是摘要。

同一個規則寫成 PySpark,就是下面這段。前一行只產生計畫,show() 才讓它執行;寫出(write)同樣是動作。

from pyspark.sql import functions as F

# df 是已存在的 DataFrame。下一行只把步驟寫進執行計畫。
planned = df.filter(df.amount > 100).groupBy("region").agg(F.sum("amount"))
planned.show()  # 動作:這裡才開始計算

Apache Kafka解決的是來源與處理系統綁死的問題。它是分散式訊息佇列(message queue),也是一份可保留的事件日誌(event log)。生產者把事件寫入主題(topic),主題分成多個分區(partition),同一個消費者群組(consumer group)裡的消費者並行讀取不同分區。處理端可以事後重播日誌,不必請來源再送一次。Kafka 負責把事件送到、並把底稿留住;分數、視窗與模型不是它的工作。

Apache Flink負責計算這些事件,而且是原生的逐筆串流。時間可以用事件時間(event time),也就是事件發生的時間;處理時間(processing time)則是引擎收到它的時間。網路會讓兩者岔開。視窗決定一段事件如何收成一個結果。滾動視窗(tumbling window)首尾相接、互不重疊,例如每 5 分鐘一塊。滑動視窗(sliding window)以較短的間隔向前移,相鄰視窗可以重疊,同一筆事件會落入一個以上的視窗。工作階段視窗(session window)不看固定時鐘,而用閒置間隔切開:中間空了一段沒有新事件,上一個工作階段就結束。水位線(watermark)是一條隨著事件前進的時間界線,用來處理晚到的事件,判斷它還能不能補進已經打開的視窗。

兩種架構決定批次與串流要不要各留一層。Lambda 架構讓批次層與速度層並行:批次層算得完整但較慢,速度層補上最新的一段,使用時把兩邊的結果合併。Kappa 架構只留串流層;邏輯要重算時,重播事件日誌再跑一次。

進倉儲的順序也有兩種。ETL(Extract, Transform, Load)先轉換再載入,是傳統資料倉儲的做法。ELT(Extract, Load, Transform)先載入,再在目標系統裡轉換;雲端資料倉儲與湖倉(lakehouse)常見這種順序,因為轉換可以靠目標端的算力來做。

步驟之間有先後、還要定時重跑時,Apache Airflow把資料管線(data pipeline)編成 DAG 來排程。這張 DAG 管的是「哪一步在哪一步之後」,和 Spark 那張「這一步裡的轉換如何執行」的計畫不是同一層。

程式若留在 Python,而不走 Spark,還有兩條路。Dask用接近 pandas 與 NumPy 的介面,把表格與陣列運算拆到多核心或多台機器。Ray面向分散式 Python 與 AI 工作負載,例如分散式訓練(distributed training),以及平行做超參數搜尋(hyperparameter search)。

局部合併可以用標準差對一次帳。母體變異數寫成 Σx²/n − (Σx/n)²,分母是 n。每個分區只要交出三個數:筆數 n、總和 Σx、平方和 Σx²。合併時三個數各自相加,再套進同一條公式,得到全體變異數,開平方就是標準差。10 億筆明細不必集中到一台機器。各分區的標準差不能直接平均:分區的平均若不相同,那個平均會丟掉分區之間的差異。第 6 題用兩小筆假資料把差額算出來。

三、容易混淆的地方

實務上容易把兩種都能處理大資料的工具混在一起。先看資料是累積後再算,還是一到就算;再看中間結果在磁碟還是記憶體;最後看工具是在送事件,還是在算事件。

容易混淆 差別在哪 實際遇到時的線索
批次處理 vs 串流處理 批次累積一段再算,延遲高、吞吐大;串流資料一到就算,延遲在秒級以下 報表每天更新時用批次處理;即時服務需在 1 秒內判斷時用串流處理
MapReduce vs Spark 都能分散式計算;MapReduce 中間結果寫磁碟,多輪迭代慢;Spark 留在記憶體,迭代型工作較快 工作紀錄顯示機器學習需迭代數十輪且重視速度時,Spark 較合適
Kafka vs Flink Kafka 傳輸並暫存事件,日誌可重播;Flink 逐筆計算,處理視窗、事件時間與水位線 架構需求提到主題、分區、解耦與重播時對應 Kafka;要處理視窗、晚到資料或即時評分時對應 Flink
轉換 vs 動作 轉換只建立執行計畫;動作才觸發執行 程式碼只有 filter、select、groupBy,沒有 count、collect、show、write,執行很快結束且沒有結果,表示只有轉換、尚未觸發動作
Lambda vs Kappa Lambda 是批次層加速層,結果再合併;Kappa 只留串流層,重算靠重播事件日誌 架構文件描述批次層與加速層並行、結果再合併時是 Lambda;只保留串流層並靠重播日誌重算時是 Kappa
ETL vs ELT ETL 先轉換再載入;ELT 先載入,在目標系統內轉換 資料流程在進傳統倉儲前清洗時偏 ETL;在雲端倉儲或湖倉載入後轉換時偏 ELT
邊緣運算 vs 雲端 邊緣在產生端附近做低延遲前處理,原始資料可留在現場;雲端做集中訓練與跨點彙整 工廠現場頻寬有限且要求毫秒級處理時用邊緣;需要集中訓練或跨點彙整時用雲端

四、基礎練習

第 1 題

某銀行的信用卡授權規定:每筆交易都要在 1 秒內判斷是不是盜刷,逾時就只能先放行。交易在尖峰時仍持續進來。哪一種做法對得上這個時限?

  • (A) 每天凌晨用 MapReduce 重算前一日交易,隔天把高風險卡號列入名單
  • (B) 每週把交易匯出到資料倉儲,由分析師在週一產出風險名單
  • (C) 用串流處理接住交易:Kafka 接收事件,Flink 在事件到達時計算特徵並評分
  • (D) 值班人員每小時抽查一批交易,確認異常後再人工鎖卡
看答案與解析

答案:C

判斷必須在 1 秒內完成,資料一到就要算,這是串流處理的時間尺度。Kafka 接收並暫存事件,Flink 在到達時計算特徵與評分,分工對得上授權當下的決策。

  • (A) 每日批次要到隔天才能產出名單,超過 1 秒的限制。
  • (B) 每週匯出更慢,倉儲名單也不是授權當下的判斷。
  • (D) 人工抽查涵蓋不到每一筆,也無法在 1 秒內做完。

第 2 題

某電商要在全部瀏覽與購買紀錄上訓練模型,同一份大資料必須反覆掃描 50 輪。團隊在 MapReduce 與 Spark 之間做選擇。哪一個判斷站得住?

  • (A) MapReduce 每一輪都把中間結果寫入磁碟,所以會比留在記憶體更快
  • (B) 兩者都是分散式運算,迭代 50 輪的耗時沒有差別
  • (C) 迭代訓練應交給 Kafka,因為消費者群組可以並行讀取分區
  • (D) Spark 把中間結果留在記憶體,50 輪不必每輪讀寫磁碟,因此比 MapReduce 快
看答案與解析

答案:D

這是迭代型工作。Spark 的中間結果在記憶體,MapReduce 每一輪把中間結果寫入磁碟再讀出。輪數到 50,差別來自磁碟往返,而不是兩者能不能算。

  • (A) 反覆寫磁碟是 MapReduce 在多輪訓練裡變慢的原因。
  • (B) 中間結果放在磁碟或記憶體,多輪之後耗時不同。
  • (C) Kafka 負責事件的傳輸與暫存,不負責模型的迭代訓練。

第 3 題

工程師在 PySpark 寫了 filter 與 groupBy,後面沒有 show()、count()、collect() 或 write。程式很快結束,也沒有任何彙總結果。原因是什麼?

  • (A) 惰性求值下,這些轉換只建立執行計畫,要等動作才真正計算
  • (B) 資料已經在叢集上算完,結果被引擎自動丟棄,所以看不到
  • (C) Spark 計算出錯,但惰性求值把錯誤藏起來,因此沒有報錯
  • (D) groupBy 是動作,算完才會結束;結束得快代表資料是空的
看答案與解析

答案:A

filter 與 groupBy 都是轉換,只把步驟加進執行計畫。沒有 count、collect、show、write 這類動作,引擎不會去讀資料,所以程式很快結束,彙總也還沒產生。

  • (B) 計畫尚未執行,不存在已經算完的結果。
  • (C) 工作還沒送出,沒有這一次計算的錯誤可藏;快結束是因為計畫未執行。
  • (D) groupBy 是轉換,不是動作;結束得快不能用來推斷資料是空的。

第 4 題

某工廠位在頻寬有限的廠區。產線攝影機要在毫秒級判斷零件有沒有瑕疵,判斷完才決定這個零件往哪一條輸送帶走。哪一種架構合理?

  • (A) 每一幀影像即時上傳雲端,等雲端模型判完再把指令送回產線
  • (B) 在現場的閘道器做邊緣推論,只把判別結果或摘要上傳
  • (C) 影像先載入資料倉儲,用倉儲裡的查詢判斷瑕疵
  • (D) 白天只錄影,每晚用批次作業統一判定瑕疵
看答案與解析

答案:B

時限是毫秒級,而且廠區頻寬有限。推論要放在資料產生端附近:邊緣節點在現場算出結果,上傳的是結果或摘要,原始影像不必全數離開工廠。雲端仍可另收摘要,做集中訓練與跨線彙整,但那一步補不上輸送帶當下的時限。

  • (A) 整幀上傳與往返都吃頻寬,也達不到毫秒級。
  • (C) 資料倉儲的載入與查詢是分析路徑,不是產線當下的分流。
  • (D) 每日批次把判斷延到晚上,零件通過攝影機時無法分流。

第 5 題

某新聞 App 要統計每 5 分鐘的點擊數。部分手機因網路不穩,事件會晚幾分鐘才送達。若晚到的點擊仍要算進它真正發生的那 5 分鐘,應該怎麼做?

  • (A) 以處理時間開 5 分鐘視窗,事件抵達引擎的時間就是它的時間
  • (B) 晚到的事件全部丟棄;只晚幾分鐘,每 5 分鐘的點擊數仍然正確
  • (C) 取消 5 分鐘視窗,改算整段期間的全距(range),晚到就不會影響結果
  • (D) 以事件時間開視窗,並用水位線處理晚到、決定它能否補進對應視窗
看答案與解析

答案:D

點擊要歸到它發生的那 5 分鐘,所以視窗按事件時間切。水位線用來處理晚幾分鐘才到的事件,讓它還能補進正確的視窗,而不是跟著抵達時刻走。

  • (A) 處理時間會把晚到的點擊算進抵達的那一段,5 分鐘的歸屬會錯。
  • (B) 丟棄會少算該視窗的次數,正確性會受影響。
  • (C) 全距是最大值減最小值,不是點擊數,也沒有把晚到事件送回正確的 5 分鐘。

第 6 題

10 億筆資料散在多個分區,要算全體的標準差,而且不能把明細集中到一台機器。哪一種做法是對的?

  • (A) 呼叫 collect() 把 10 億筆拉回 driver,在單機計算標準差
  • (B) 只取第一個分區的標準差,用它代表全體
  • (C) 各分區先算 n、Σx、Σx²,合併這三個量後再套變異數公式
  • (D) 各分區各自算出標準差,再把這些標準差做算術平均
看答案與解析

答案:C

母體變異數是 Σx²/n − (Σx/n)²。n、Σx、Σx² 都可以跨分區相加,所以各分區交這三個數就能合併,10 億筆不必集中。

用兩小筆假資料看(D)為什麼不行。分區甲是 1 與 3:n = 2,Σx = 1 + 3 = 4,Σx² = 1 + 9 = 10。變異數 = 10/2 − (4/2)² = 5 − 4 = 1,標準差 = √1 = 1。分區乙是 10 與 12:n = 2,Σx = 22,Σx² = 100 + 144 = 244。變異數 = 244/2 − (22/2)² = 122 − 121 = 1,標準差也是 1。兩者的算術平均是 (1 + 1) / 2 = 1。

合併後 n = 2 + 2 = 4,Σx = 4 + 22 = 26,Σx² = 10 + 244 = 254。變異數 = 254/4 − (26/4)² = 63.5 − 6.5² = 63.5 − 42.25 = 21.25,標準差是 √21.25。兩區的平均分別是 2 與 11,直接平均標準差會把這段差距丟掉。

  • (A) collect() 把全部明細拉回 driver,10 億筆會讓單機記憶體不足。
  • (B) 第一個分區只描述自己的分佈,代替不了其餘分區。
  • (D) 如上,算術平均得到 1,合併公式得到的變異數是 21.25,兩者不相等。

五、重點回顧

  • 批次處理累積一段資料再算,延遲高、吞吐大,配每日報表;串流處理資料一到就算,延遲在秒級以下,配 1 秒內要完成的判斷。
  • Hadoop 用 HDFS 存放區塊,預設每塊複製 3 份以容錯,用 MapReduce 計算並把中間結果寫入磁碟,用 YARN 分配資源;多輪迭代會被反覆讀寫磁碟拖慢。
  • Spark 把中間結果留在記憶體,所以 50 輪這類迭代訓練通常比 MapReduce 快;轉換只建立執行計畫,動作才執行,collect() 會把全量拉回 driver。
  • Kafka 把事件寫入主題的分區,讓消費者群組並行讀取,並保留可重播的日誌;Flink 逐筆計算,用事件時間、滾動/滑動/工作階段視窗與水位線處理晚到。
  • Lambda 架構是批次層加速層再合併結果;Kappa 架構只留串流層,重算時重播事件日誌。ETL 先轉換再載入,ELT 先載入再在目標內轉換。頻寬有限又要毫秒級時,邊緣運算在現場處理,只上傳結果或摘要。
  • 母體變異數是 Σx²/n − (Σx/n)²。各分區交出 n、Σx、Σx² 再合併,10 億筆不必集中;各分區標準差直接平均並不正確,因為分區平均不同的那一段差異會被丟掉。

覺得有幫助? RSS · X · 請我喝杯咖啡

先看最左的拆分節點,再沿三條箭頭看到 Map、Shuffle 與 Reduce。

局部運算後再依鍵彙整
1 / 4
資料拆成區塊並留在各節點
Map 在各區塊局部處理,Shuffle 依鍵重分組,Reduce 合併同鍵結果。
  1. 運算要送到資料所在的那一台,而不是把全部明細先搬去同一個地方,這個原則稱為資料局部性。
  2. 各台機器只處理自己的區塊,讀取、排序與彙總才不會一起卡在同一台機器上。
  3. 這一步之後,網路上移動的是可相加、可合併的摘要,原始明細仍留在原地。
  4. 這裡只合併已經算好的局部結果,全部明細不必集中到同一台機器。

表格五列由上到下各是一種工具,中欄區分批次或串流,右欄寫定位。

工具分工在儲存、計算與傳送
1 / 5
HDFS、批次資料、分散式儲存
HDFS 存分散檔案,MapReduce 與 Spark 負責計算,Kafka 傳送事件,Flink 處理串流。
  1. 檔案切成區塊後散放多台機器,預設每個區塊複製 3 份,一份讀不到還有副本。
  2. 中間結果寫入磁碟。只跑一輪時影響不大;同一份資料要反覆掃很多輪,寫磁碟的時間就累積起來。
  3. 中間結果留在記憶體,迭代型工作不必每輪寫入磁碟,多輪訓練通常因此比 MapReduce 快。微批次把很短區間內到達的資料收成一小批,執行時仍是一小批一小批。
  4. 生產者把事件寫入主題後可以重播,處理端不必請來源再送一次。分數、視窗與模型不是它的工作。
  5. 網路會讓事件時間與處理時間岔開;水位線用來判斷晚到的事件還能不能補進已開的視窗。