上一章談的是「拆分與串接佇列」的模式。複製並產生多個不同輸出固然是批次處理的重要一環,但同樣重要的另一半是:把多路輸出重新拉回來,產生某種聚合輸出

  • 這種聚合最經典的例子,就是 MapReduce 中的 reduce 階段:不難看出 map 步驟正是「分片一個工作佇列」,而 reduce 步驟則是「協調式處理」——最終把大量輸出收斂成單一的聚合回應。
  • 但批次處理的聚合模式不只 reduce 一種。

圖 12-1:通用的平行工作分配與結果聚合批次系統

匯流(Join,或稱屏障同步)#

  • 先前章節看過如何把工作拆開、平行分散到多個節點(尤其是分片工作佇列)。但處理工作流時,有時必須先湊齊完整的工作集合,才能進入下一階段
  • 上一章的**合併器(merger)**是一個選項,但它只是把兩個佇列的輸出摻在一起變成一個佇列供後續處理——它並不保證處理開始前資料集是完整的。這代表你既無法保證處理的完整性,也沒有機會為所有已處理元素計算聚合統計量。
  • 我們需要更強的協調原語(primitive),那就是匯流模式(join pattern)。它類似於「join 一個執行緒」:所有工作仍平行進行,但在所有平行處理的工作項目都完成之前,任何項目都不會被放出 join。這在並行程式設計中一般稱為屏障同步(barrier synchronization)

匯流的價值與代價是一體兩面:

  • 價值:確保在進行聚合階段(例如對集合中某個值求總和)之前,資料一筆不缺
  • 代價:要求所有資料都被前一階段處理完,後續運算才能開始。這降低了工作流可能的平行度,因而拉長了整體延遲

圖 12-2:批次處理的匯流模式

歸約(Reduce)#

  • 如果分片工作佇列是 map/reduce 演算法中的 map 階段,那麼剩下的就是 reduce 階段。歸約之所以屬於協調式批次處理模式,是因為它不管輸入怎麼切分都能進行,而且用途與 join 類似:把多個不同批次操作的平行輸出組合在一起。
  • 但與匯流關鍵不同的是:歸約的目標不是等所有資料處理完,而是樂觀地(optimistically)把平行的資料項目逐步合併成對整個集合的單一綜合表示
  • 歸約的每一步把數個輸出合併成一個輸出。之所以叫「reduce」,是因為它同時做了兩件縮減:減少輸出的總數量,以及把完整資料項縮減為「回答特定批次運算所必需的代表性資料」
  • 由於歸約階段吃的是一段範圍的輸入、吐出的是相似形態的輸出,它可以依需要重複執行任意多次,直到整個資料集收斂成單一輸出。

這正是歸約相對匯流的幸運之處:歸約可以在 map/分片階段還在跑的時候就平行啟動。當然,要產出完整輸出,所有資料終究都得被處理過;但「可以提早開始」讓整體批次運算跑得更快。

動手做:計數(Count)

任務:計算一本書中某些單字出現的次數。

先用分片把計數工作切成多個工作佇列——例如建立 10 個分片佇列、由 10 個人各自負責一個佇列,依頁碼末位數分片:頁碼結尾為 1 的都給第一個佇列、結尾為 2 的給第二個,依此類推。

每個人數完自己負責的頁面後,把結果寫在紙上:

a: 50
the: 17
cat: 2
airplane: 1
...

這就可以輸出給歸約階段。歸約的作用是把兩個以上的輸出合併成單一輸出。給定第二份輸出:

a: 30
the: 25
dog: 4
airplane: 2
...

歸約的做法是把各單字的計數加總起來:

a: 80
the: 42
dog: 4
cat: 2
airplane: 3
...

顯然這個歸約階段可以在前一輪歸約的輸出上反覆執行,直到只剩單一歸約輸出——這也意味著歸約本身可以平行進行。最終的輸出,就是整本書所有單字的完整計數。

求和(Sum)#

  • 與計數相似但略有不同的歸約形式:把一群不同的值加總起來。差別在於不是每個值都算「1」,而是把原始輸出資料中實際帶的數值相加。
  • 例子——統計美國總人口:先量測每個鎮的人口,再全部加總。
    • 第一步:把工作依分片成多個工作佇列。這是很好的第一層分片,但即使平行處理,單一個人要數完一州所有鎮的人口仍然太久。
    • 因此再做第二層分片,這次依切分。至此我們先平行化到州、再平行化到郡,每個郡的佇列產出一串 (鎮, 人口) 的元組。
  • 此時歸約模式就能上場,而且它甚至不需要知道我們做了兩層分片:歸約只要抓兩個以上的輸出項(例如 (Seattle, 4,000,000)(Northampton, 25,000))相加,產生新輸出 (Seattle-Northampton, 4,025,000) 即可。
  • 如同計數,這個歸約可以用同一份程式碼重複執行任意次數,最後只剩一個包含全美總人口的輸出——而且幾乎所有運算都是平行發生的。

直方圖(Histogram)#

  • 歸約模式的最後一個例子:在用平行分片/map 與歸約統計美國人口的同時,我們也想建立「美國平均家庭」的模型——具體來說,是一張家庭規模的直方圖,估計有零到十個小孩的家庭各有多少。多層分片的做法完全照舊(甚至可以沿用同一批工作者)。

  • 這次資料蒐集階段的輸出,是每個鎮一張直方圖

    0: 15%
    1: 25%
    2: 50%
    3: 10%
    4: 5%
  • 乍看之下,直方圖似乎很難合併。但只要結合求和例子裡的人口資料就迎刃而解:

    • 把每張直方圖乘上它對應的相對人口,即可得到每個項目的總人數
    • 再把這個新總數除以合併後的人口總和,就能把多張不同的直方圖合併更新為單一輸出。
  • 於是同樣可以反覆套用歸約模式,直到產生單一輸出為止。

動手做:影像標記與處理管線

任務:對一大批尖峰時段的公路影像做標記與處理——計算汽車、卡車、機車的數量,以及各汽車顏色的分布;此外還有一個前置步驟,要把所有車牌模糊化以保護隱私。影像以一系列 HTTPS URL 的形式送達,每個 URL 指向一張原始影像。

階段一:找出並模糊車牌。 為簡化佇列中的每項任務,這裡用兩個工作者——一個偵測車牌位置、一個把該位置模糊化——並以上一章的多工作者模式把兩者組成單一容器群組。這種關注點分離看似多餘,但很有價值:模糊化工作者可以被重用來模糊其他輸出(例如人臉)。同時,為了確保可靠性並最大化平行處理,我們把影像分片到多個工作者佇列。

圖 12-3:分片工作佇列與多個模糊化分片

階段二:匯流後刪除原圖。 每張影像成功模糊後會上傳到另一個位置,接著才刪除原圖。但在所有影像都成功模糊之前,我們不想刪除原圖——萬一發生災難性故障需要重跑整條管線就慘了。因此用匯流模式把所有分片模糊佇列的輸出合併成單一佇列,只有在所有分片都完成後才釋出項目。

階段三:複製器分岔。 現在可以刪除原圖,並開始車款與顏色偵測。為了最大化管線吞吐量,用複製器模式把工作項目複製到兩個佇列:

  • 一個負責刪除原始影像。
  • 一個負責辨識車輛類型(汽車、卡車、機車)與顏色。

圖 12-4:管線中的輸出匯流、複製器、刪除與影像辨識部分

階段四:分片辨識與歸約聚合。 辨識車輛與顏色的佇列同樣先套用分片模式分散到多個佇列,每個佇列有兩個工作者:一個辨識每輛車的位置與類型、一個辨識某區域的顏色,再以多工作者模式接合。同樣地,把程式碼分到不同容器,讓顏色偵測容器能被重用到辨識車色以外的多種任務上。

該佇列的輸出是這樣的 JSON 元組(代表單一影像中發現的資訊):

{
  "vehicles": {
    "car": 12,
    "truck": 7,
    "motorcycle": 4
  },
  "colors": {
    "white": 8,
    "black": 3,
    "blue": 6,
    "red": 6
  }
}

最後,用前述由 MapReduce 發揚光大的歸約模式把所有資料加總在一起——做法與計數例子完全相同——產出整批影像的最終影像數與顏色計數。

本章小結#

  • 協調式批次處理處理的是「把平行輸出重新聚合」這一半問題。
  • **匯流(join)**提供屏障同步:保證資料完整,代價是犧牲平行度、拉長延遲。
  • **歸約(reduce)**樂觀地兩兩合併,可反覆套用直到單一輸出;因為能與 map 階段平行進行,整體運算更快。計數、求和、直方圖都是它的實例。
  • 真實管線(如影像標記)通常是分片、多工作者、匯流、複製器、歸約多種模式的組合。