11. 批處理

帶有太強個人色彩的系統無法成功。當最初的設計完成並且相對穩健時,不同的人開始以自己的方式進行試驗,真正的考驗才開始。
高德納
到目前為止,本書的大部分內容都在討論 請求 和 查詢,以及相應的 響應 或 結果。許多現代資料系統都預設採用這種處理方式:你請求某樣東西,或者發出一條指令,系統便儘量快速地給出答案。
瀏覽器請求網頁、服務呼叫遠端 API,以及資料庫、快取、搜尋索引等許多系統,都是這樣工作的。我們稱它們為 線上系統。這類系統通常以響應時間作為主要效能指標,而且往往需要具備容錯能力,才能保證高可用性。
然而,有些計算規模太大,或者需要處理的資料太多,無法放在一次互動式請求中完成。例如,你可能需要訓練 AI 模型,把大量資料從一種形式轉換成另一種形式,或者在非常大的資料集上進行分析。我們把這類任務稱為 批處理 作業,相應的系統有時也稱為 離線系統。
批處理作業讀取一組輸入資料(只讀),併產生一組輸出資料(每次執行都從頭生成)。它通常不會像讀寫事務那樣修改現有資料。因此,輸出是由輸入 衍生 而來的(參見“權威記錄系統與衍生資料”):如果對輸出不滿意,只需把它刪除,調整作業邏輯,再執行一次。把輸入視為不可變資料,並避免產生副作用(例如寫入外部資料庫),不僅能讓批處理作業獲得良好的效能,還會帶來其他好處:
如果程式碼中引入了錯誤,導致輸出有誤或遭到破壞,只要回滾到先前版本的程式碼並重新執行作業,輸出便能恢復正確。更簡單的辦法是把舊輸出儲存在另一個目錄中,需要時直接切換回去。大多數物件儲存和開放表格式(參見“雲資料倉庫”)都支援這種稱為 時間旅行 的功能。大多數支援讀寫事務的資料庫卻不具備這種性質:如果錯誤程式碼把壞資料寫進資料庫,回滾程式碼並不能修復已經寫入的資料。這種從錯誤程式碼中恢復的能力稱為 容忍人為失誤 1。
由於回滾很容易,功能開發可以比“犯錯就會造成不可逆損害”的環境推進得更快。這種 儘量減少不可逆操作 的原則有利於敏捷軟體開發 2。
同一組檔案可以作為許多不同作業的輸入,其中也包括監控作業:它們計算指標,並檢查某項作業的輸出是否具備預期特徵,例如將其與上一次執行的輸出比較,衡量兩者之間的差異。
批處理框架能夠高效利用計算資源。雖然 OLTP 資料庫和應用伺服器等線上資料系統也能成批處理資料,但完成同樣工作所需的資源可能昂貴得多。
批處理也會帶來一些挑戰。在大多數框架中,只有整個作業執行完畢,其他作業才能繼續處理它的輸出。批處理還可能效率不高:輸入資料發生任何變化——哪怕只有一個位元組——都意味著批處理作業必須重新處理整個輸入資料集。儘管存在這些侷限,批處理仍在許多場景中證明了自己的價值,我們將在“批處理用例”中再次討論這些場景。
一項批處理作業可能要執行很長時間:幾分鐘、幾小時,甚至幾天。作業也可能按固定週期排程執行,例如每天一次。其主要效能指標通常是吞吐量,即單位時間內能夠處理多少資料。有些批處理系統遇到故障時只會中止並重新啟動整個作業;另一些則具備容錯能力,即使某些節點崩潰,作業也能順利完成。
Note
批處理之外還有一種選擇,即 流處理。流處理作業不會在處理完當前輸入後結束,而是繼續監視輸入,並在輸入發生變化後不久加以處理。我們將在第 12 章討論流處理。
線上系統與批處理系統之間的界限並不總是分明:一條長時間執行的資料庫查詢看起來就很像批處理過程。不過,批處理有一些獨特的性質,使其成為構建可靠、可伸縮且可維護應用的重要構件。例如,它經常用於 資料整合,也就是把多個數據系統組合起來,完成單個系統無法獨自完成的工作。“資料倉庫”中討論的 ETL 就是一例。
現代批處理深受 MapReduce 影響。Google 於 2004 年發表了這種批處理演算法 3,隨後 Hadoop、CouchDB 和 MongoDB 等多個開源資料系統都實現了它。MapReduce 是一種相當底層的程式設計模型,沒有資料倉庫中的並行查詢執行引擎那麼精巧 4 5。剛問世時,MapReduce 讓普通商用硬體能夠達到的處理規模邁上了一個新臺階;如今它已經基本過時,Google 也不再使用 6 7。
如今,批處理更多由 Spark、Flink 之類的框架或資料倉庫查詢引擎完成。與 MapReduce 一樣,這些系統高度依賴分片(參見第 7 章)和並行執行,但它們的快取和執行策略精巧得多。隨著系統逐漸成熟,運維方面的問題基本得到解決,關注點也轉向了易用性。資料流 API、查詢語言和資料框 API 如今都得到了廣泛支援,作業與工作流編排同樣日趨成熟。Oozie、Azkaban 等以 Hadoop 為中心的工作流排程器,已經被 Airflow、Dagster 和 Prefect 等更通用的方案所取代;後者支援各種批處理框架和雲資料倉庫。
雲計算已經十分普及,批處理的儲存層也正從 HDFS、GlusterFS 和 CephFS 等分散式檔案系統(DFS)轉向 S3 之類的物件儲存。BigQuery、Snowflake 等可伸縮的雲資料倉庫,則進一步模糊了資料倉庫與批處理系統之間的界限。
為了直觀地理解批處理,本章先從單臺機器上的標準 Unix 工具入手,再研究如何把資料處理擴充套件到分散式系統中的多臺機器。我們將看到,分散式批處理框架與作業系統很相似,同樣擁有排程器和檔案系統。隨後,我們會考察編寫批處理作業時使用的幾種處理模型,最後討論常見的批處理用例。
使用 Unix 工具的批處理
假設有一臺 Web 伺服器,每處理一個請求,它都會在日誌檔案末尾追加一行。以 nginx 預設的訪問日誌格式為例,其中一行可能如下所示:
216.58.210.78 - - [27/Jun/2025:17:55:11 +0000] "GET /css/typography.css HTTP/1.1"
200 3377 "https://martin.kleppmann.com/" "Mozilla/5.0 (Macintosh; Intel Mac OS X
10_15_7) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/137.0.0.0 Safari/537.36"
(這其實是一行,只是為了便於閱讀才在這裡折成了多行。)這一行包含許多資訊。要解釋它,需要先檢視日誌格式的定義:
$remote_addr - $remote_user [$time_local] "$request"
$status $body_bytes_sent "$http_referer" "$http_user_agent"
這行日誌表示:在 UTC 時間 2025 年 6 月 27 日 17:55:11,伺服器收到了來自客戶端 IP 地址 216.58.210.78、請求檔案 /css/typography.css 的請求。使用者沒有經過身份認證,因此 $remote_user 被設為連字元(-)。響應狀態碼為 200(即請求成功),響應大小為 3,377 位元組。瀏覽器是 Chrome 137;它之所以載入這個檔案,是因為網址 https://martin.kleppmann.com/ 的頁面引用了該檔案。
解析日誌聽起來或許像是一個刻意編造的例子,實際上卻是許多現代科技公司的關鍵工作,從廣告資料管道到支付處理無所不包。事實上,日誌處理正是 MapReduce 得以迅速普及、並推動“大資料”浪潮的重要原因之一。
簡單日誌分析
許多工具都能讀取這些日誌檔案,生成漂亮的網站流量報告。不過為了練習,我們來用基本的 Unix 工具自行構建一個。假設你想找出網站上最受歡迎的五個頁面,可以在 Unix shell 中執行下面的命令:
cat /var/log/nginx/access.log | #1
awk '{print $7}' | #2
sort | #3
uniq -c | #4
sort -r -n | #5
head -n 5 #6讀取日誌檔案。(嚴格來說,這裡的
cat並非必需,因為可以把輸入檔案直接作為引數傳給awk。不過這樣寫能讓線性管道顯得更清楚。)按空白字元把每一行拆成欄位,只輸出第七個欄位,而它恰好就是請求的 URL。在前面的示例中,這個 URL 是 /css/typography.css。
按字母順序
sort請求 URL 列表。如果某個 URL 被請求了 n 次,那麼排序後的檔案中就會連續出現 n 個相同的 URL。uniq命令透過檢查相鄰兩行是否相同,濾掉輸入中重複的行。-c選項讓它同時輸出計數器:對於每個不同的 URL,它會報告該 URL 在輸入中出現了多少次。第二個
sort按每行開頭的數字(-n)排序,也就是按 URL 的請求次數排序;然後以逆序(-r)返回結果,讓最大的數字排在最前面。最後,
head只輸出輸入的前五行(-n 5),並丟棄其餘內容。
這一系列命令的輸出大致如下:
4189 /favicon.ico
3631 /2016/02/08/how-to-do-distributed-locking.html
2124 /2020/11/18/distributed-systems-and-elliptic-curves.html
1369 /
915 /css/typography.css
如果你不熟悉 Unix 工具,前面的命令列可能顯得有些晦澀,但它的能力非常強大。它能在幾秒鐘內處理數 GB 的日誌,而且可以輕鬆修改分析方式來滿足需要。例如,如果不希望報告中包含 CSS 檔案,只需把 awk 引數改為 '$7 !~ /\.css$/ {print $7}';如果想統計最常出現的客戶端 IP 地址,而不是最常訪問的頁面,則把引數改為 '{print $1}',以此類推。
本書沒有足夠篇幅詳細介紹 Unix 工具,但它們非常值得學習。令人驚訝的是,僅用 awk、sed、grep、sort、uniq 和 xargs 的某種組合,幾分鐘內就能完成許多資料分析,而且效能也出奇地好 8。
命令鏈與自定義程式
除了使用 Unix 命令鏈,你也可以寫一個簡單的程式來完成同樣的工作。例如,Python 程式可能如下所示:
from collections import defaultdict
counts = defaultdict(int) #1
with open('/var/log/nginx/access.log', 'r') as file:
for line in file:
url = line.split()[6] #2
counts[url] += 1 #3
top5 = sorted(((count, url) for url, count in counts.items()), reverse=True)[:5] #4
for count, url in top5: #5
print(f"{count} {url}")counts是一個散列表,儲存每個 URL 出現次數的計數器;計數器的預設值為 0。從每一行日誌中取出按空白字元分隔的第七個欄位,作為 URL(因為 Python 陣列從 0 開始計數,所以陣列下標是 6)。
把當前日誌行中 URL 對應的計數器加一。
按計數器的值對散列表內容降序排列,並取出前五項。
輸出這五項。
這個程式沒有 Unix 管道那麼簡潔,但也相當容易讀懂;喜歡哪一種,部分取決於個人偏好。不過,除了表面的語法差異,兩者的執行流程也大不相同。如果在大檔案上執行這項分析,區別就會顯現出來。
排序與記憶體聚合
Python 指令碼在記憶體中維護一個 URL 散列表,把每個 URL 對映到它出現的次數。Unix 管道沒有這樣的散列表,而是依靠排序 URL 列表;同一個 URL 出現多次時,它只是在列表中重複多次。
哪種方法更好?這取決於不同 URL 的數量。對於大多數中小型網站,你大概可以把所有不同的 URL 及其計數器放進 1 GB 左右的記憶體。這個作業的 工作集(即作業需要隨機訪問的記憶體量)只取決於不同 URL 的數量:即使日誌中同一個 URL 出現了一百萬次,散列表所需的空間仍然只是一個 URL 加一個計數器。只要工作集足夠小,記憶體散列表就能很好地工作——即便在膝上型電腦上也是如此。
另一方面,如果作業的工作集大於可用記憶體,排序方法就有一個優勢:它可以高效利用磁碟。這與“日誌結構儲存”中討論的原理相同:先在記憶體中對資料塊排序,並把它們作為段檔案寫入磁碟,再將多個有序段合併成一個更大的有序檔案。歸併排序採用順序訪問模式,在磁碟上表現很好(參見“順序與隨機寫入”)。
GNU Coreutils(Linux)中的 sort 工具會把無法裝入記憶體的資料自動溢寫到磁碟,還會自動利用多個 CPU 核並行排序 9。這意味著前面的簡單 Unix 命令鏈可以輕鬆擴充套件到大型資料集,而不會耗盡記憶體。此時的瓶頸很可能是從磁碟讀取輸入檔案的速度。
Unix 工具的侷限在於,它們只能在單臺機器上執行。如果資料集大到無法裝入單機記憶體或本地磁碟,就會遇到問題——這正是分散式批處理框架的用武之地。
分散式系統中的批處理
執行前面 Unix 工具示例的機器,需要由幾個元件協同處理日誌資料:
透過作業系統的檔案系統介面訪問的儲存裝置。
決定程序何時執行、如何為其分配 CPU 資源的排程器。
一系列 Unix 程式,其
stdin和stdout透過管道連線在一起。
分散式資料處理框架中也存在同樣的元件。事實上,你可以把分散式處理框架看作一種分散式作業系統:它們擁有檔案系統和作業排程器,程式則透過檔案系統或其他通訊通道相互發送資料。
分散式檔案系統
作業系統提供的檔案系統由若干層組成:
最底層的塊裝置驅動程式直接與磁碟通訊,讓上層能夠讀寫原始資料塊。
塊裝置層之上是頁快取,它把最近訪問的資料塊儲存在記憶體中,以加快再次訪問的速度。
檔案系統層封裝了塊 API,把大檔案拆分成塊,並維護 inode、目錄和檔案等元資料。例如,ext4 和 XFS 是 Linux 上常用的兩種實現。
最後,作業系統透過名為虛擬檔案系統(VFS)的統一 API,把不同檔案系統暴露給應用。無論底層採用哪一種檔案系統,應用都可以用同樣的方式讀寫資料。
分散式檔案系統的工作方式與此非常相似。檔案同樣被拆分成塊,只不過這些塊分佈在許多機器上。分散式檔案系統的塊通常比本地檔案系統大得多:HDFS(Hadoop 分散式檔案系統)的預設塊大小為 128 MB,JuiceFS 和許多物件儲存使用 4 MB 的塊,而 ext4 的塊只有 4,096 位元組。塊越大,需要跟蹤的元資料就越少;對於 PB 級資料集,這種差異十分顯著。相對於讀取資料塊所需的時間,大塊也能降低尋道開銷所佔的比例。
大多數物理儲存裝置無法寫入不完整的資料塊,因此即使資料沒有填滿一個塊,作業系統也必須為寫入使用整個塊。分散式檔案系統的塊更大,而且通常構建在作業系統檔案系統之上,所以沒有這種要求。例如,一個 900 MB 的檔案採用 128 MB 的分塊時,會由 7 個佔用 128 MB 的塊和 1 個佔用 4 MB 的塊組成。
讀取分散式檔案系統中的塊,需要向叢集中存放該塊的機器發出網路請求。每臺機器都執行一個守護程序,並公開一套 API,讓遠端程序能夠讀寫以檔案形式儲存在其本地檔案系統中的資料塊。HDFS 把這些守護程序稱為 DataNode,GlusterFS 則稱為 glusterfsd。本書統一把它們稱為 資料節點。
分散式檔案系統還實現了分散式版本的頁快取。由於資料塊以檔案形式存放在資料節點上,讀寫操作會經過每個資料節點的作業系統,其中就包含記憶體頁快取。因此,經常讀取的資料塊會被快取在資料節點的記憶體中。有些分散式檔案系統還實現了更多快取層,例如 JuiceFS 提供的客戶端快取和本地磁碟快取。
ext4 和 XFS 等檔案系統會跟蹤空閒空間、檔案塊位置、目錄結構和許可權設定等儲存元資料。分散式檔案系統同樣需要記錄檔案分散在哪些機器上、具有什麼許可權等資訊。Hadoop 透過 NameNode 服務維護叢集元資料;DeepSeek 的 3FS 則使用元資料服務,並把資料持久化到 FoundationDB 之類的鍵值儲存中。
檔案系統層之上是 VFS;在批處理中,與之最接近的是分散式檔案系統的協議。分散式檔案系統必須公開某種協議或介面,讓批處理系統能夠讀寫檔案。它充當一套可插拔介面:只要實現該協議,任何分散式檔案系統都可以接入。例如,MinIO、Cloudflare R2、Tigris 和 Backblaze B2 等許多儲存系統都採用了 Amazon S3 API。支援 S3 的批處理系統,也就能夠使用其中任何一種儲存。
有些分散式檔案系統提供與 POSIX 相容的介面,在作業系統的 VFS 看來,它們與其他檔案系統並無不同。這類系統通常透過使用者空間檔案系統(FUSE)或網路檔案系統(NFS)協議接入 VFS。NFS 或許是最著名的分散式檔案系統協議。它最初的用途,是讓多個客戶端讀寫單臺伺服器上的資料。近來,AWS Elastic File System(EFS)和 Archil 等檔案系統提供了可伸縮得多、但仍與 NFS 相容的分散式實現。NFS 客戶端依舊只連線一個端點,不過這些系統會在內部與分散式元資料服務和資料節點通訊,完成資料讀寫。
分散式檔案系統與網路儲存
分散式檔案系統以 無共享 原則為基礎(參見“共享記憶體、共享磁碟與無共享架構”),不同於網路附加儲存(NAS)和儲存區域網路(SAN)架構採用的 共享磁碟 方法。共享磁碟儲存由集中式儲存裝置實現,往往需要定製硬體和光纖通道等專用網路基礎設施。相比之下,無共享方法不需要特殊硬體,只需用普通的資料中心網路連線計算機即可。
許多分散式檔案系統構建在普通商用硬體上。這種硬體比較便宜,但故障率也高於企業級硬體。為了容忍機器和磁碟故障,檔案塊會複製到多臺機器上。這樣也便於排程器更均勻地分配工作負載,因為任務可以在任何存有其輸入資料副本的節點上執行。這裡的複製可以像第 6 章所述,在多臺機器上儲存若干份完整副本;也可以採用 Reed–Solomon 碼等 糾刪碼,以低於完整複製的儲存開銷恢復丟失的資料 10 11 12。這些技術與 RAID 很相似,後者在連線到同一臺機器的多個磁碟之間提供冗餘;區別在於,分散式檔案系統透過普通的資料中心網路訪問和複製檔案,不需要特殊硬體。
物件儲存
Amazon S3、Google Cloud Storage、Azure Blob Storage 和 OpenStack Swift 等物件儲存服務,已經成為批處理作業中分散式檔案系統的常用替代方案。事實上,兩者的界限有些模糊。正如上一節和“以物件儲存為後端的資料庫”中所述,使用者空間檔案系統(FUSE)驅動程式可以讓使用者把 S3 之類的物件儲存當作檔案系統使用。JuiceFS 和 Ceph 等分散式檔案系統實現也同時提供物件儲存與檔案系統 API。不過,兩類系統的 API、效能和一致性保證可能大相徑庭。即使某個系統看似實現了所需的 API,採用之前也必須仔細確認它的實際行為符合預期。
物件儲存中的每個物件都有一個 URL,例如 s3://my-photo-bucket/2025/04/01/birthday.png。URL 的主機部分(my-photo-bucket)表示存放物件的儲存桶(bucket),後面的部分則是物件的 鍵(本例中為 /2025/04/01/birthday.png)。儲存桶的名稱在全域性範圍內唯一,而每個物件的鍵在所屬儲存桶內必須唯一。
物件透過 get 呼叫讀取,透過 put 呼叫寫入。與檔案系統中的檔案不同,物件寫入後是不可變的;若要更新,只能像鍵值儲存一樣,透過 put 完整重寫整個物件。Azure Blob Storage 和 S3 Express One Zone 支援追加寫入,但大多數物件儲存都不支援。物件儲存也沒有 fopen、fseek 之類的檔案控制代碼 API。
物件看起來似乎按目錄組織,但這有些誤導,因為物件儲存根本沒有目錄的概念。路徑結構只是一種約定,斜槓本身也是物件鍵的一部分。這種約定允許你按特定字首請求物件列表,實現類似目錄列表的效果。不過,按字首列出物件與檔案系統列出目錄有兩點區別:
字首
list操作類似於 Unix 系統上的遞迴ls -R:它會返回所有以該字首開頭的物件,其中也包括子路徑下的物件。物件儲存中不可能存在空目錄。假如刪除
s3://my-photo-bucket/2025/04/01下的所有物件,那麼對s3://my-photo-bucket/2025/04呼叫list時,01就不會再出現。常見的做法是用一個零位元組物件來表示空目錄,例如建立空物件s3://my-photo-bucket/2025/04/01,這樣即使所有子物件都已刪除,它仍會保留下來。
分散式檔案系統通常支援硬連結、符號連結、檔案鎖和原子重新命名等常見檔案操作,物件儲存則沒有這些功能。它們一般不支援連結和鎖,重新命名也不是原子的,而是先把物件複製到新鍵,再刪除舊物件。如果想重新命名一個“目錄”,就必須逐一重新命名其中的每個物件,因為目錄名本身是物件鍵的一部分。
第 4 章討論的鍵值儲存針對小值(通常只有幾 KB)以及頻繁、低延遲的讀寫進行了最佳化。相比之下,分散式檔案系統和物件儲存通常針對大物件(數 MB 至數 GB)和頻率較低、規模較大的讀取進行了最佳化。不過最近,物件儲存也開始支援更加頻繁、規模更小的讀寫。例如,S3 Express One Zone 如今可以達到個位數毫秒延遲,其定價模型也更接近鍵值儲存。
分散式檔案系統與物件儲存還有一個區別:HDFS 等分散式檔案系統可以把計算任務安排在存有特定檔案副本的機器上執行。任務就能直接讀取本地檔案,無需透過網路傳輸;如果任務的可執行程式碼遠小於它要讀取的檔案,這樣可以節省大量頻寬。物件儲存通常把儲存與計算分開。這種做法可能消耗更多頻寬,但現代資料中心網路的速度很快,因而往往可以接受。儲存與計算解耦後,CPU、記憶體等計算資源還可以獨立於儲存容量進行伸縮。
分散式作業編排
作業系統的類比同樣適用於作業編排。執行 Unix 批處理作業時,總要有某個元件真正執行 awk、sort、uniq 和 head 這些程序。它必須把一個程序的輸出傳到另一個程序的輸入,為每個程序分配記憶體,在 CPU 上公平地排程並執行各程序的指令,實施記憶體和 I/O 邊界,等等。在單臺機器上,這些工作由作業系統核心負責;在分散式環境中,它們則由作業編排器承擔。
批處理框架會向編排器的排程器傳送執行作業的請求。啟動作業的請求包含如下元資料:
要執行的任務數量;
每項任務所需的記憶體、CPU 和磁碟資源;
作業識別符號;
訪問憑據;
輸入、輸出資料等作業引數;
GPU 或磁碟型別等必要的硬體資訊;
作業可執行程式碼所在的位置。
Kubernetes 和 Hadoop YARN(Yet Another Resource Negotiator)13 等編排器會結合這些資訊與叢集元資料,透過下列元件執行作業:
- 任務執行器(Task Executor)
叢集中的每個節點都執行著執行器守護程序,例如 YARN 的 NodeManager 或 Kubernetes 的 kubelet。執行器負責執行作業任務,傳送心跳來表明自己仍然存活,並跟蹤節點上的任務狀態與資源分配。執行器收到啟動任務的請求後,會取得作業的可執行程式碼,並執行命令來啟動任務。然後,它會監視程序,直到程序結束或失敗,再相應更新任務狀態元資料。
許多執行器還會與作業系統配合,同時提供安全隔離與效能隔離。例如,YARN 和 Kubernetes 都使用 Linux cgroups。這樣既能防止任務訪問無權訪問的資料,也能防止任務過度使用資源、影響同一節點上其他任務的效能。
- 資源管理器(Resource Manager)
編排器的資源管理器儲存每個節點的元資料,包括可用硬體(CPU、GPU、記憶體和磁碟等)、任務狀態、網路位置、節點狀態以及其他相關資訊。因此,資源管理器能夠提供叢集當前狀態的全域性檢視。資源管理器的中心化特性可能同時形成可伸縮性和可用性瓶頸。YARN 使用 ZooKeeper、Kubernetes 使用 etcd 來儲存叢集狀態(參見“協調服務”)。
- 排程器(Scheduler)
編排器通常有一個中心化的排程子系統,負責接收啟動、停止作業或查詢作業狀態的請求。例如,排程器可能收到一項請求:使用特定的 Docker 映象,在配有某種 GPU 的節點上啟動一項包含 10 個任務的作業。排程器根據請求中的資訊和資源管理器儲存的狀態,決定把哪些任務放在哪些節點上執行。隨後,它會把分配結果通知任務執行器,由執行器開始執行任務。
雖然各個編排器使用的術語不盡相同,但幾乎所有編排系統中都能找到這些元件。
Note
有些排程決策需要由應用專用的排程器來完成,以便考慮特定需求,例如在查詢量達到某個閾值時自動擴充套件只讀副本。中心排程器與應用專用排程器共同決定任務的最佳執行方式。YARN 把這種子排程器稱為 ApplicationMaster,Kubernetes 則稱為 operator。
資源分配
排程器在作業編排中扮演著格外棘手的角色:面對需求相互競爭的作業,它必須找出分配叢集有限資源的最佳方式。從根本上說,排程決策需要在公平與效率之間取得平衡。
設想一個由五個節點組成的小型叢集,總共有 160 個 CPU 核可用。叢集排程器收到了兩項作業請求,每項都希望使用 100 個核來完成工作。怎樣排程才最好?
排程器可以同時為每項作業執行 80 個任務,等先前的任務完成後,再啟動兩項作業各自剩餘的 20 個任務。
排程器也可以先執行一項作業的所有任務,等到有 100 個核可用時,再開始執行第二項作業。這種策略稱為 成組排程(gang scheduling)。
一項作業請求會比另一項先到。排程器必須決定是把 100 個核全部分配給先到的作業,還是留下一些資源,為尚未到來的作業做準備。
這個例子非常簡單,卻已經暴露出許多艱難的權衡。以成組排程為例:如果排程器不斷預留 CPU 核,直到 100 個核能夠同時使用,那麼一些節點就會閒置,叢集的資源利用率也會下降;如果其他作業也試圖預留 CPU 核,甚至還可能發生死鎖。
另一方面,如果排程器只是等待 100 個核空閒,其他作業又可能在此期間搶走這些核。叢集或許會在很長時間內都湊不出 100 個可用核,從而導致 飢餓。排程器還可以 搶佔 第一項作業的部分任務,終止它們來為第二項作業騰出空間。不過,搶佔任務同樣會降低叢集效率,因為被終止的任務稍後需要重新啟動、重新執行。
現在再設想一下,排程器必須為數百乃至數百萬項這樣的作業請求作出分配決策。要找到最優解似乎根本不可行。事實上,這個問題是 NP-hard 的;也就是說,除了規模最小的例子之外,計算最優解所需的時間長得令人無法接受 14 15。
因此,實際的排程器會採用啟發式方法,作出雖非最優、但還算合理的決策。常見演算法包括先進先出(FIFO)、主導資源公平(DRF)、優先順序佇列、基於容量或配額的排程,以及各種裝箱演算法。這些演算法的細節超出了本書範圍,但排程確實是一個十分有趣的研究領域。
工作流排程
本章開頭的 Unix 工具示例由若干命令串聯而成。分散式批處理也經常採用相同的模式:一項作業的輸出需要成為另一項或多項作業的輸入,而一項作業又可能有多個輸入,分別由其他作業產生。這種作業結構稱為 工作流,也稱為作業的 有向無環圖(DAG)。
Note
在“持久化執行與工作流”中,我們見過能夠持久執行一系列步驟的工作流引擎,這些步驟通常會發出 RPC。在批處理的語境中,“工作流”有不同含義:它是一系列批處理過程,每個過程都接收輸入資料、產生輸出資料,通常不會向外部服務發出 RPC。持久化執行引擎通常會比批處理系統在每個請求中處理更少的資料,不過兩者之間的界限並不十分清晰。
採用多項作業組成的工作流可能有幾個原因:
如果一項作業的輸出需要成為多項其他作業的輸入,而且這些下游作業由不同團隊維護,最好先讓第一項作業把輸出寫到一個所有下游作業都能讀取的位置。每當資料更新時,消費它的作業就可以安排執行,也可以按照其他時間表執行。
你可能需要把資料從一種處理工具傳給另一種。例如,一項 Spark 作業把資料輸出到 HDFS,隨後由 Python 指令碼觸發 Trino SQL 查詢(參見“雲資料倉庫”),繼續處理 HDFS 檔案,並把結果輸出到 S3。
有些資料管道本身就需要多個處理階段。例如,某個階段需要按一個鍵分片資料,而下一階段需要按另一個鍵分片,那麼第一個階段就可以按照第二階段需要的方式對輸出資料分片。
在 Unix 工具示例中,連線一項命令輸出與另一項命令輸入的管道只使用一個很小的記憶體緩衝區,並不會把資料寫入檔案。如果緩衝區已滿,生產資料的程序就必須等待,直到消費程序從緩衝區中讀走一部分資料,才能繼續輸出——這是一種 背壓。Spark、Flink 等批處理執行引擎支援類似的模型,可以把一項任務的輸出直接傳給另一項任務;如果兩項任務執行在不同機器上,資料就透過網路傳輸。
不過在工作流中,更常見的做法是讓一項作業把輸出寫入分散式檔案系統或物件儲存,再由下一項作業從那裡讀取。這樣可以使作業彼此解耦,在不同時間執行。如果一項作業有多個輸入,工作流排程器通常要等到產生這些輸入的所有作業都成功完成,才會執行消費這些輸入的作業。
YARN ResourceManager 等編排框架中的排程器,以及 Spark 的內建排程器,都不會管理完整的工作流;它們只按單項作業進行排程。為了處理多次作業執行之間的依賴,人們開發了 Airflow、Dagster 和 Prefect 等工作流排程器。維護大量批處理作業時,這些排程器提供了非常有用的管理功能。許多資料管道的工作流通常包含 50 至 100 項作業;在大型組織中,還可能有許多團隊執行不同的作業或工作流,跨越多個系統讀取彼此的輸出。管理這樣複雜的資料流,離不開相應的工具支援。
故障處理
批處理作業往往會執行很長時間。一項包含許多並行任務、長時間執行的作業,很可能會在途中遇到至少一次任務失敗。正如“硬體與軟體故障”和“不可靠的網路”中所討論的,這可能由許多原因造成,其中包括硬體故障(普通商用硬體上尤其常見)和網路中斷。
任務無法完成的另一個原因,是排程器可能有意搶佔(終止)它。當系統設定多個優先順序時,搶佔尤其有用:低優先順序任務執行成本較低,高優先順序任務則要付出更高成本。只要還有空閒計算容量,就可以執行低優先順序任務;但如果一項高優先順序任務到來,低優先順序任務隨時可能遭到搶佔。這類價格較低的低優先順序虛擬機器,在 Amazon EC2 中稱為 競價例項(spot instance),在 Azure 中稱為 競價虛擬機器(spot virtual machine),在 Google Cloud 中則稱為 可搶佔例項(preemptible instance)16。
批處理通常用於時效性要求不高的作業,因此很適合使用低優先順序任務和競價例項,以降低執行成本。實質上,這些作業利用了原本會被閒置的計算資源,從而提高叢集利用率。不過,這也意味著排程器更可能終止這些任務:搶佔發生的頻率要高於硬體故障 17。
由於批處理作業每次執行都從頭生成輸出,任務失敗要比線上系統中的故障容易處理:系統可以刪除失敗執行留下的不完整輸出,再把任務安排到另一臺機器上重新執行。不過,只因一項任務失敗就重跑整個作業會非常浪費。因此,MapReduce 及其後繼系統讓並行任務彼此獨立,從而可以按單項任務的粒度重試工作 3。
如果一項任務的輸出要在工作流中成為另一項任務的輸入,容錯就會更加棘手。MapReduce 的辦法是始終把這類中間資料寫回分散式檔案系統,並等待寫入任務順利完成,之後才允許其他任務讀取資料。即使在搶佔頻繁的環境中,這種做法也能正常工作,但它要向分散式檔案系統寫入大量資料,效率可能很低。
Spark 把中間資料儲存在記憶體中,必要時“溢寫”到本地磁碟,只把最終結果寫入分散式檔案系統。它還會跟蹤中間資料的計算過程,以便在資料丟失時重新計算 18。Flink 採用另一種方法,定期為任務狀態的快照建立檢查點 19。我們將在“資料流引擎”中再次討論這個話題。
批處理模型
我們已經瞭解了分散式環境如何排程批處理作業。現在把注意力轉向批處理框架實際處理資料的方式。最常見的兩種模型是 MapReduce 和資料流引擎。雖然在實踐中,資料流引擎已經基本取代 MapReduce,但理解 MapReduce 的工作原理仍然很有用,因為許多現代批處理框架都深受它的影響。
MapReduce 和資料流引擎已經逐漸支援多種程式設計模型,其中包括底層程式設計 API、關係查詢語言和資料框 API。豐富的選擇使應用工程師、分析工程師、業務分析師,甚至不具備技術背景的員工,都能夠出於各種用途處理企業資料。我們將在“批處理用例”中討論這些用途。
MapReduce
MapReduce 的資料處理模式與“簡單日誌分析”中的 Web 伺服器日誌示例非常相似:
讀取一組輸入檔案,並把它們拆分成 記錄。在 Web 伺服器日誌示例中,每條記錄就是日誌的一行(也就是說,
\n是記錄分隔符)。在 Hadoop MapReduce 中,輸入檔案儲存在 HDFS 之類的分散式檔案系統,或 S3 之類的物件儲存中。系統可以使用多種檔案格式,例如 Apache Parquet(列式格式,參見“列式儲存”)或 Apache Avro(行式格式,參見“Avro”)。對每條輸入記錄呼叫 mapper 函式,從中提取鍵和值。在 Unix 工具示例中,mapper 函式就是
awk '{print $7}':它提取 URL($7)作為鍵,並把值留空。按鍵對所有鍵值對排序。在日誌示例中,這一步由第一個
sort命令完成。呼叫 reducer 函式,遍歷排序後的鍵值對。如果同一個鍵出現多次,排序會讓這些記錄在列表中彼此相鄰,因此不必在記憶體中儲存大量狀態,就能輕鬆合併這些值。在 Unix 工具示例中,reducer 由
uniq -c命令實現,它負責統計具有相同鍵的相鄰記錄數。
這四個步驟可以由一項 MapReduce 作業完成。第 2 步(map)與第 4 步(reduce)需要編寫定製的資料處理程式碼。第 1 步(把檔案拆成記錄)由輸入格式解析器負責。第 3 步的 sort 在 MapReduce 中是隱式的——mapper 的輸出總會在傳給 reducer 之前排序,所以不需要自行實現。排序是批處理中的基礎演算法,我們將在“混洗資料”中再次討論。
要建立一項 MapReduce 作業,需要實現 mapper 和 reducer 兩個回撥函式,其行為如下:
- Mapper
每條輸入記錄都會呼叫一次 mapper,其任務是從輸入記錄中提取鍵和值。對於每條輸入,它可以生成任意數量的鍵值對(也可以一個都不生成)。mapper 不會把某條輸入記錄的狀態保留到下一條,因此每條記錄都能獨立處理。
- Reducer
MapReduce 框架接收 mapper 產生的鍵值對,收集屬於同一個鍵的所有值,再用一個遍歷該值集合的迭代器呼叫 reducer。reducer 可以生成輸出記錄,例如同一 URL 的出現次數。
在 Web 伺服器日誌示例中,第 5 步還有第二個 sort 命令,用於按請求次數排列 URL。在 MapReduce 中,如果需要第二個排序階段,可以再編寫一項 MapReduce 作業,把第一項作業的輸出作為第二項的輸入。從這個角度看,mapper 的作用是準備資料,把它轉換成適合排序的形式;reducer 的作用則是處理已經排好序的資料。
MapReduce 與函數語言程式設計
MapReduce 雖然用於批處理,其程式設計模型卻來自函數語言程式設計。Lisp 最早引入了 map 和 reduce(也稱 fold),把它們作為列表上的高階函式;後來,Python、Rust 和 Java 等主流語言也採用了這些函式。包括 SQL 提供的操作在內,許多常見的資料處理操作都可以建立在 MapReduce 之上。兩個函式乃至函數語言程式設計整體所具備的一些重要性質,恰好能為 MapReduce 所用。map 與 reduce 可以相互組合,這非常適合資料處理(正如 Unix 示例所示)。map 還天然易於並行,因為每項輸入都獨立處理;reduce 則可以並行處理不同的鍵。
實際上,使用原始 MapReduce API 實現複雜的處理作業非常困難而且費力——例如,作業使用的任何連線演算法都必須從頭實現 20。與較新的批處理器相比,MapReduce 的速度也很慢。其中一個原因是,基於檔案的 I/O 無法形成作業流水線:上游作業完成之前,下游作業不能開始處理它的輸出。
資料流引擎
為了解決 MapReduce 的一些問題,人們開發了幾種新的分散式批處理執行引擎,其中最著名的是 Spark 18 21 和 Flink 19。它們的設計方式各有不同,卻有一個共同點:把整個工作流作為一項作業處理,而不是拆成彼此獨立的子作業。
這些系統顯式建模了資料流經多個處理階段的過程,因此稱為 資料流引擎。與 MapReduce 一樣,它們提供底層 API,透過反覆呼叫使用者定義的函式,每次處理一條記錄;同時也提供 連線 和 分組 等高層運算元。它們透過對輸入分片來並行執行工作,並把一項任務的輸出複製到另一項任務,作為後者的輸入;如果兩項任務執行在不同機器上,複製就透過網路完成。與 MapReduce 不同,運算元不必嚴格交替扮演 map 和 reduce 的角色,而可以以更加靈活的方式組合。
這些資料流 API 通常使用關係模型風格的構件來表達計算:按照某個欄位的值連線資料集,按鍵對元組分組,根據某項條件過濾資料,以及透過計數、求和等函式聚合元組。這些操作在內部透過下一節討論的混洗演算法實現。
這類處理引擎建立在 Dryad 22 和 Nephele 23 等研究系統的基礎之上。與 MapReduce 模型相比,它們有幾個優點:
排序之類代價高昂的工作只需在確有必要之處執行,不必默認出現在每個 map 階段與 reduce 階段之間。
如果一連串運算元都不改變資料集的分片方式(例如 map 或 filter),它們就可以合併成一項任務,從而減少複製資料的開銷。
工作流中的所有連線和資料依賴都經過顯式宣告,因此排程器可以總覽全域性,瞭解各處需要哪些資料,並據此最佳化資料區域性。例如,它可以嘗試把消費某項資料的任務放在產生該資料的任務所在機器上,這樣便能透過共享記憶體緩衝區交換資料,無需透過網路複製。
運算元之間的中間狀態通常只需儲存在記憶體中,或寫入本地磁碟;這樣所需的 I/O 少於寫入分散式檔案系統或物件儲存,因為後者還要把資料複製到多臺機器上,並在每個副本所在機器上寫入磁碟。MapReduce 已經對 mapper 的輸出採用了這種最佳化,資料流引擎則把這個思想推廣到了所有中間狀態。
運算元可以在輸入準備就緒後立即開始執行,無需等前一階段全部結束,下一階段才能開始。
現有程序可以重複用於執行新的運算元,從而減少啟動開銷;相比之下,MapReduce 會為每項任務啟動一個新的 JVM。
資料流引擎能夠實現與 MapReduce 工作流相同的計算,而且得益於上述最佳化,執行速度通常快得多。
混洗資料
我們已經看到,本章開頭的 Unix 工具示例和 MapReduce 都以排序為基礎。批處理器必須能夠對 PB 級的資料集排序,而這些資料不可能裝入一臺機器。因此,它們需要一種輸入和輸出都經過分片的分散式排序演算法,這種演算法稱為 混洗(shuffle)。
混洗不是隨機
“shuffle” 這個詞容易引起誤解:把一副撲克牌洗牌,會得到隨機順序;這裡所說的混洗卻會產生排好序的結果,其中不含任何隨機性。
混洗是批處理器的一項基礎演算法,連線與聚合都要使用它。MapReduce、Spark、Flink、Daft、Dataflow 和 BigQuery 24 都實現了可伸縮的高效能混洗演算法,以處理大型資料集。下面以 Hadoop MapReduce 中的混洗為例 25,不過本節的概念也適用於其他系統。
圖 11-1 展示了一項 MapReduce 作業中的資料流。假設作業的輸入已經分片,三個分片分別標為 m 1、m 2 和 m 3。例如,每個分片可以是 HDFS 中的單獨檔案,也可以是物件儲存中的單獨物件。同一資料集的全部分片,可以集中放在同一個 HDFS 目錄中;在物件儲存的儲存桶裡,它們也可以使用相同的鍵字首。

框架會為每個輸入分片啟動一項單獨的 map 任務。任務讀取分配給它的檔案,每次把一條記錄傳給 mapper 回撥函式。計算的 reduce 端同樣經過分片。map 任務的數量由輸入分片數決定,而 reduce 任務的數量則由作業作者配置,兩者可以不同。
mapper 的輸出由鍵值對組成。框架必須保證:如果兩個不同的 mapper 輸出了相同的鍵,這些鍵值對最終會由同一個 reducer 任務處理。為此,每個 mapper 都會在本地磁碟上為每個 reducer 分別建立一個輸出檔案。例如,圖 11-1中的檔案 m 1, r 2 由 mapper 1 建立,包含發往 reducer 2 的資料。mapper 輸出一個鍵值對時,通常會對鍵進行雜湊,據此決定把它寫入哪個 reducer 檔案,這與“按鍵的雜湊分片”相似。
mapper 在寫入這些檔案的同時,還會在每個檔案中按鍵排列鍵值對。這裡可以採用“日誌結構儲存”中見過的技術:先在記憶體的有序資料結構中收集一批鍵值對,再把它們寫成有序段檔案,然後逐步將較小的段合併成較大的段。
每個 mapper 完成後,reducer 會連線到它,並把屬於自己的有序鍵值對檔案複製到本地磁碟。reduce 任務取得所有 mapper 輸出中屬於自己的那一份後,會像歸併排序一樣合併這些檔案,同時保持排序順序。這樣一來,即使具有相同鍵的鍵值對來自不同的 mapper,最終也會彼此相鄰。隨後,每個鍵都會呼叫一次 reducer 函式,並向它傳入一個迭代器,用於返回該鍵對應的所有值。
reducer 函式產生的記錄會順序寫入檔案,每項 reduce 任務對應一個檔案。圖 11-1中的 r 1、r 2 和 r 3 就構成了作業輸出資料集的三個分片,它們會被寫回分散式檔案系統或物件儲存。
MapReduce 在 map 與 reduce 階段之間執行混洗,而現代資料流引擎和雲資料倉庫則要精巧得多。BigQuery 等系統優化了混洗演算法,儘量把資料儲存在記憶體中,或把資料寫入外部排序服務 24。這類服務既能加快混洗,又能透過複製混洗資料來提高容錯能力。
JOIN 與 GROUP BY
下面來看有序資料如何簡化分散式連線與聚合。為便於說明,我們仍以 MapReduce 為例,不過這些概念適用於大多數批處理系統。
圖 11-2 展示了批處理作業中一個典型的連線示例。左側是一份事件日誌,記錄已登入使用者在網站上的操作,這些記錄稱為 活動事件,也叫 點選流資料;右側則是使用者資料庫。這個例子可以看作星型模式的一部分(參見“星型與雪花型:分析模式”):事件日誌是事實表,使用者資料庫則是其中一張維度表。

如果想結合使用者資料庫中的資訊來分析活動事件——例如利用使用者資料中的出生日期,瞭解某些頁面更受年輕使用者還是年長使用者歡迎——就需要連線這兩張表。假設兩張表都大到必須分片,該如何計算這種連線?
可以利用 MapReduce 的一個性質:無論鍵值對最初位於哪個分片,混洗都會把具有相同鍵的鍵值對集中到同一個 reducer。這裡可以把使用者 ID 作為鍵。因此,我們可以編寫一個 mapper 遍歷使用者活動事件,以使用者 ID 為鍵,輸出頁面瀏覽的 URL,如圖 11-3所示。另一個 mapper 逐行遍歷使用者資料庫,提取使用者 ID 作為鍵、使用者出生日期作為值。

混洗隨後確保 reducer 函式能夠同時訪問某位使用者的出生日期,以及該使用者的所有頁面瀏覽事件。MapReduce 甚至可以安排記錄的順序,讓 reducer 總是先看到使用者資料庫中的記錄,接著再按時間戳順序看到活動事件;這種技術稱為 二次排序 25。
這樣,reducer 就能輕鬆完成實際的連線邏輯。第一個值應當是出生日期,reducer 先把它儲存在區域性變數中,再遍歷具有相同使用者 ID 的活動事件,輸出每個瀏覽過的 URL 以及瀏覽者的出生日期。reducer 會一次處理某個使用者 ID 的全部記錄,所以任何時刻只需在記憶體中儲存一條使用者記錄,也完全不必發出網路請求。這種演算法稱為 排序合併連線(sort-merge join),因為 mapper 的輸出按鍵排序,而 reducer 隨後會合並連線兩側的有序記錄列表。
工作流中的下一項 MapReduce 作業可以繼續計算每個 URL 的瀏覽者年齡分佈。它首先用 URL 作為鍵混洗資料;排序完成後,reducer 會遍歷同一個 URL 的所有頁面瀏覽記錄(其中包含瀏覽者的出生日期),為各個年齡段維護瀏覽次數計數器,併為每條頁面瀏覽記錄遞增相應的計數器。這樣就實現了 分組(group by)操作和聚合。
查詢語言
多年來,分散式批處理的執行引擎日趨成熟。如今,基礎設施已經足夠穩健,可以在超過 10,000 臺機器的叢集上儲存和處理數 PB 資料。既然如此規模的批處理系統如何實際執行已經大體得到解決,人們就把注意力轉向了程式設計模型的改進。
MapReduce、資料流引擎和雲資料倉庫都採用 SQL 作為批處理的通用語言。這是很自然的選擇:傳統資料倉庫本來就使用 SQL,資料分析和 ETL 工具也已經支援 SQL,而且開發者和分析師無不熟悉 SQL。
與手寫 MapReduce 作業相比,查詢語言介面除了所需程式碼更少,還有一個明顯的優點:它們可以互動使用。你可以編寫分析查詢,再從終端或圖形介面執行。這種互動查詢方式很高效也很自然,業務分析師、產品經理、銷售和財務團隊等人員都可以藉此在批處理環境中探索資料。儘管它不是經典形式的批處理,SQL 支援仍然使探索性查詢成為分散式批處理系統的一項適用場景。
高階查詢語言不僅提高了使用系統的人所能達到的生產率,也能從機器層面提高作業執行效率。正如“雲資料倉庫”所述,查詢引擎負責把 SQL 查詢轉換成在叢集中執行的批處理作業。從查詢轉換到語法樹、再轉換成物理運算元的過程,讓引擎有機會最佳化查詢。Hive、Trino、Spark 和 Flink 等查詢引擎都提供基於代價的查詢最佳化器,可以分析連線輸入的特徵,自動決定哪一種演算法最適合眼前的任務。最佳化器甚至可以改變連線的順序,以儘量減少中間狀態 19 26 27 28。
SQL 是最流行的通用批處理查詢語言,不過其他語言仍用於一些專門場景。Apache Pig 是一種以關係運算元為基礎的語言,允許使用者逐步描述資料管道,而不是把所有邏輯寫成一個龐大的 SQL 查詢。資料框(見下一節)具有類似特徵,Morel 則是受 Pig 影響的一種較新語言。還有一些使用者採用 jq、JMESPath 或 JsonPath 等 JSON 查詢語言。
在“圖資料模型”中,我們討論了如何用圖來建模資料,以及如何用圖查詢語言遍歷圖中的邊和頂點。許多圖處理框架也支援透過查詢語言進行批計算,例如 Apache TinkerPop 的 Gremlin。我們將在“批處理用例”中進一步討論圖處理的用例。
批處理與雲資料倉庫正在收斂
過去,資料倉庫執行在專用硬體裝置上,為關係資料提供 SQL 分析查詢。相比之下,MapReduce 等批處理框架的目標是提供更強的可伸縮性與靈活性:它們支援使用通用程式語言編寫處理邏輯,並允許讀寫任意資料格式。
隨著時間推移,兩者變得越來越相似。現代批處理框架如今支援用 SQL 編寫批處理作業,並透過 Parquet 等列式儲存格式和經過最佳化的查詢執行引擎,在關係查詢上取得了良好效能(參見“查詢執行:編譯與向量化”)。與此同時,資料倉庫遷移到雲端後,變得更加可伸縮(參見“雲資料倉庫”),並實現了許多與分散式批處理框架相同的排程、容錯和混洗技術。其中許多系統也使用分散式檔案系統。
正如批處理系統採用 SQL 作為處理模型一樣,雲資料倉庫也採用了資料框等其他處理模型(見下一節)。例如,Google Cloud BigQuery 提供 BigQuery DataFrames 庫,Snowflake 的 Snowpark 則與 pandas 整合。Airflow、Prefect 和 Dagster 等批處理工作流編排器同樣能夠整合雲資料倉庫。
當然,並非所有批處理作業都容易用 SQL 表達。PageRank 等迭代圖演算法、複雜的機器學習以及許多其他任務,都很難用 SQL 編寫。AI 資料處理包含影象、影片和音訊等非關係型多模態資料,同樣很難用 SQL 完成。
此外,雲資料倉庫不擅長某些工作負載。使用列式儲存格式時,逐行計算的效率較低;在這種情況下,最好使用資料倉庫的其他 API,或者改用批處理系統。雲資料倉庫也往往比其他批處理系統昂貴。大型作業改在 Spark 或 Flink 等批處理系統上執行,可能更加經濟。
歸根結底,應該用批處理系統還是資料倉庫來處理資料,取決於成本、便利性、實現難度、可用性等因素。大多數大型企業都擁有多套資料處理系統,因而可以靈活作出選擇;小型公司則往往只靠一套系統就已足夠。
資料框
隨著資料科學家和統計學家開始使用分散式批處理框架進行機器學習,他們發現現有的處理模型用起來很麻煩,因為他們習慣的是 R 和 pandas 中的資料框模型(參見“資料框、矩陣與陣列”)。資料框與關係資料庫中的表相似:它由許多行組成,同一列中的所有值都具有相同型別。使用者不必編寫一條龐大的 SQL 查詢,而是呼叫與關係運算元對應的函式,執行過濾、連線、排序、分組等操作。
最初,資料框操作一般在本機記憶體中執行,因而只能處理單臺機器能夠容納的資料集。資料科學家希望繼續使用熟悉的資料框 API,與批處理環境中的大型資料集互動。Spark、Flink 和 Daft 等分散式資料處理框架為滿足這種需求,也採用了資料框 API。不過,本地資料框通常帶有索引並且有序,分散式資料框一般卻不具備這些性質 29。因此,把程式遷移到批處理框架後,其效能可能出人意料。
資料框 API 看起來與資料流 API 相似,但具體實現各不相同。pandas 會在資料框方法被呼叫時立即執行操作;Apache Spark 則會先把所有資料框 API 呼叫轉換成查詢計劃,經過查詢最佳化,再在分散式資料流引擎上執行工作流。這樣便能改善效能。
Daft 等框架甚至同時支援客戶端與服務端計算:規模較小的記憶體操作在客戶端執行,規模較大的資料集和計算則在服務端執行。Apache Arrow 等列式儲存格式提供了統一的資料模型,可由客戶端與服務端的執行引擎共享。
批處理用例
瞭解批處理如何工作之後,我們來看看它在各種應用中怎樣發揮作用。批處理作業非常適合成批處理大型資料集,卻不適合低延遲場景。因此,只要資料量很大、資料新鮮度又不重要,通常就能看到批處理作業。聽起來這似乎是一個很大的侷限,但事實證明,相當多的資料處理都符合這一模型:
會計與庫存核對通常成批進行,企業藉此檢查交易是否與銀行賬戶和庫存相符 30。
製造業的需求預測由週期執行的批處理作業計算 31。
許多金融系統同樣以批處理為基礎。例如,美國的銀行網路幾乎完全依靠批處理作業執行 34。
下面幾節將討論幾乎每個行業都能見到的一些批處理用例。
提取—轉換—載入(ETL)
“資料倉庫”介紹了 ETL 和 ELT:資料處理管道從生產資料庫抽取資料,對其進行轉換,再把結果載入到下游系統。本節用“ETL”同時指代 ETL 與 ELT 工作負載。這類工作負載經常由批處理作業完成,尤其是在下游系統為資料倉庫時。
批處理作業天然具有並行性,因而非常適合資料轉換。許多資料轉換工作負載都易於並行:過濾資料、投影欄位,以及其他許多常見的資料倉庫轉換,都可以並行完成。
批處理環境還配有穩健的工作流排程器,能夠輕鬆安排、編排和除錯 ETL 資料管道作業。發生故障時,排程器通常會重試作業,以消除可能出現的暫時性問題。作業若反覆失敗,就會被標記為失敗,開發者可以很容易地看到資料管道中的哪項作業停止了工作。Airflow 等排程器甚至內建了 MySQL、PostgreSQL、Snowflake、Spark 和 Flink 等數十種流行系統的資料來源、資料接收端及查詢運算元。排程器與資料處理系統的緊密整合簡化了資料整合。
我們還看到,批處理作業出了問題之後很容易排查和修復;這一點在除錯資料管道時尤其有用。可以直接檢查有問題的檔案,找出錯誤所在;修復 ETL 批處理作業後,再重新執行即可。例如,輸入檔案可能不再包含轉換作業所要使用的某個欄位。資料工程師發現欄位缺失後,可以更新轉換邏輯,或者修改產生該輸入的作業。
過去,資料管道通常由一支資料工程團隊統一管理,因為要求開發產品功能的其他團隊編寫和管理複雜的批處理資料管道並不公平。近來,批處理模型和元資料管理的改進,使組織中的工程師更容易參與和管理自己的資料管道。資料網格(data mesh)35 36、資料契約(data contract)37 和 資料編織(data fabric)38 等實踐提供了標準與工具,幫助團隊安全地釋出資料,供組織中的任何人使用。
如今,資料管道與分析查詢不僅開始共享處理模型,也開始共享執行引擎。許多批處理 ETL 作業與讀取其輸出的分析查詢,會執行在同一套系統上。資料管道轉換和分析查詢都以 SparkSQL、Trino 或 DuckDB 查詢來執行,已經十分常見。這樣的架構進一步模糊了應用工程、資料工程、分析工程和業務分析之間的界限。
分析
在“分析型與事務型系統”中,我們看到分析查詢(OLAP)經常掃描大量記錄,並執行分組與聚合。這樣的工作負載可以與其他批處理工作負載一起,在批處理系統中執行。分析師編寫 SQL 查詢,由查詢引擎執行,並讀寫分散式檔案系統或物件儲存。表與檔案之間的對映、名稱和型別等表元資料,則透過 Apache Iceberg 等表格式和 Unity 等目錄服務管理(參見“雲資料倉庫”)。這種架構稱為 資料湖倉 39。
與 ETL 一樣,SQL 查詢介面的改進使許多組織如今也使用 Spark 等批處理框架進行分析。這類查詢模式分為兩種:
預聚合查詢:把資料彙總成 OLAP 多維資料集或資料集市,以加快查詢(參見“物化檢視與多維資料集”)。預聚合資料可以在資料倉庫中查詢,也可以推送到 Apache Druid 或 Apache Pinot 等專用的實時 OLAP 系統。預聚合通常按固定週期進行,這類工作負載由“工作流排程”中討論的工作流排程器管理。
即席查詢(ad hoc query):使用者執行查詢來回答特定業務問題、調查使用者行為、除錯執行故障,以及完成其他許多工作。在這種場景中,響應時間很重要。分析師會反覆執行查詢,在得到響應、進一步瞭解正在研究的資料後,繼續調整查詢。能夠快速執行查詢的批處理框架,可以減少分析師的等待時間。
SQL 支援還讓批處理框架能夠與電子表格和 Tableau、Power BI、Looker、Apache Superset 等資料視覺化工具整合。例如,Tableau 提供 SparkSQL 和 Presto 聯結器;Apache Superset 則支援 Trino、Hive、Spark SQL、Presto 等許多最終會執行批處理作業來查詢資料的系統。
機器學習
機器學習(ML)經常用到批處理。資料科學家、機器學習工程師和 AI 工程師使用批處理框架來探索資料規律、轉換資料並訓練機器學習模型。常見用途包括:
- 特徵工程:對原始資料進行過濾和轉換,使其能夠用於訓練模型。預測模型通常要求輸入數值資料,因此工程師必須把文字或離散值等其他形式的資料轉換成所需格式。
- 模型訓練:訓練資料是批處理過程的輸入,訓練所得的模型權重則是輸出。
- 批次推理:訓練好的模型可以對大量資料成批進行預測,適用於資料集很大而又不要求實時返回結果的情況,其中也包括在測試資料集上評估模型的預測效果。
批處理框架為這些場景提供了專門的工具。例如,Apache Spark 的 MLlib 和 Apache Flink 的 FlinkML 都內建了豐富的特徵工程工具、統計函式和分類器。
推薦引擎、排序系統等機器學習應用也大量使用圖處理(參見“圖資料模型”)。許多圖演算法可以表述為:每次沿一條邊遍歷,將一個頂點與相鄰頂點連線起來以傳播某些資訊,並不斷重複,直到滿足某項條件——例如沒有更多的邊可以繼續遍歷,或者某項指標已經收斂。
批次同步並行(bulk synchronous parallel,BSP)計算模型 40 已經成為批次處理圖資料時的常用模型。Apache Giraph 20、Spark 的 GraphX API 和 Flink 的 Gelly API 41 等系統都實現了這一模型。它也稱為 Pregel 模型,因為 Google 的 Pregel 論文推廣了這種圖處理方法 42。
批處理也是大語言模型(LLM)資料準備與訓練的重要組成部分。網站等原始文字輸入通常存放在分散式檔案系統或物件儲存中,必須先經過預處理才能用於訓練。適合交給批處理框架完成的預處理步驟包括:
- 從 HTML 中提取純文字,並修復格式損壞的文字;
- 檢測並刪除質量低劣、內容無關或重複的文件;
- 對文字進行分詞(拆分成單詞),再將其轉換成嵌入向量,也就是每個單詞的數值表示。
Kubeflow、Flyte 和 Ray 等批處理框架就是為這類工作負載設計的。例如,OpenAI 在 ChatGPT 的訓練過程中就使用了 Ray 43。這些框架內建了與 PyTorch、TensorFlow、XGBoost 等大語言模型和 AI 庫的整合,並直接支援特徵工程、模型訓練、批次推理和微調(針對特定用例調整基礎模型)等操作。
最後,資料科學家還經常在 Jupyter 或 Hex 等互動式筆記本中對資料進行實驗。筆記本由若干 單元格(cell)組成,每個單元格都是一小段 Markdown、Python 或 SQL。按順序執行這些單元格,可以生成電子表格、圖表或資料。許多筆記本透過資料框 API 使用批處理系統,或者用 SQL 查詢這類系統。
對外提供衍生資料
批處理作業經常用於構建預計算或衍生資料集,例如商品推薦、面向使用者的報表,以及機器學習模型所需的特徵。這些資料集通常由生產資料庫、鍵值儲存或搜尋引擎對外提供。不論使用哪種系統,預計算資料最終都必須從批處理系統的分散式檔案系統或物件儲存,回到承載線上流量的資料庫中。
最直觀的選擇或許是直接在批處理作業中呼叫資料庫客戶端庫,逐條寫入資料庫伺服器。這種方法確實能用——前提是防火牆規則允許批處理環境直接訪問生產資料庫——但基於以下幾個原因,它並不是個好主意:
- 每條記錄都發起一次網路請求,要比批處理任務的正常吞吐量慢幾個數量級。即使客戶端庫支援批次寫入,效能也很可能不理想。
- 批處理框架通常會並行執行許多工。如果所有任務都以批處理應有的速率同時寫入同一個輸出資料庫,資料庫很容易不堪重負,查詢效能也可能隨之下降,繼而給系統其他部分帶來執行故障 44。
- 批處理作業通常會對作業輸出提供乾淨利落的“全有或全無”保證:如果作業成功,那麼即使途中有些任務失敗並經過重試,結果也等同於每個任務都恰好一次執行所產生的輸出;如果整個作業失敗,則不會產生任何輸出。然而,從作業內部寫入外部系統會產生無法用這種方式隱藏的外部可見副作用。這樣一來,就不得不考慮尚未完整的作業結果被其他系統看到的問題;任務失敗並重新啟動時,也可能重複寫入失敗執行已經產生的輸出。
更好的方案是讓批處理作業把預計算資料集推送到 Kafka 主題之類的流中,我們將在第 12 章中進一步討論這種做法。Elasticsearch 等搜尋引擎、Apache Pinot 和 Apache Druid 等實時 OLAP 系統、Venice 等衍生資料儲存 45,以及 ClickHouse 等雲資料倉庫,都內建了從 Kafka 攝取資料的能力。讓資料經過流系統,可以緩解上述問題中的一部分:
- 流系統針對順序寫入進行了最佳化,因此更適合批處理作業的大批次寫入負載;
- 流系統還可以在批處理作業與生產資料庫之間充當緩衝區。下游系統可以限制自己的讀取速率,確保仍有足夠能力承載線上流量;
- 同一個批處理作業的輸出可以由多個下游系統消費;
- 流系統可以充當批處理環境與生產網路之間的安全邊界:它可以部署在所謂的 DMZ(隔離區)網路中,位於批處理網路與生產網路之間。
不過,讓資料經過流並不會自動解決前面提到的“全有或全無”保證問題。為此,批處理作業必須在完成時通知下游系統:作業已經完成,資料現在可以對外提供。流消費者則必須能夠在收到完成通知之前,讓接收到的資料對查詢保持不可見,就像採用 讀已提交(read committed)隔離級別時未提交的事務一樣(參見“讀已提交”)。
另一種模式在資料庫初始化時更為常見:直接在批處理作業 內部 構建一個全新的資料庫,再把分散式檔案系統、物件儲存或本地檔案系統中的檔案批次載入到該資料庫中。許多資料系統都提供了相應的批次匯入工具,例如 TiDB 的 Lightning 工具,以及 Apache Pinot 和 Apache Druid 的 Hadoop 匯入作業。RocksDB 也提供了從批處理作業批次匯入 SST 的 API。
透過批處理構建資料庫並批次匯入資料,速度非常快,也更便於系統在不同資料集版本之間進行原子切換。另一方面,由批處理作業構建全新資料庫時,要增量更新資料集會比較困難。如果同時需要初始化和增量載入,通常會採用混合方案。例如,Venice 支援混合儲存,既可以進行基於行的批次更新,也可以切換整個資料集。
本章小結
本章探討了批處理系統的設計與實現。我們先從經典的 Unix 工具鏈(awk、sort、uniq 等)入手,以此說明排序、計數等基本的批處理原語。
接著,我們把規模擴大到分散式批處理系統。我們看到,批處理式 I/O 以不可變、有界的輸入資料集為處理物件並生成輸出資料,因而可以在不產生副作用的情況下重跑和除錯。為了處理檔案,批處理框架主要由三部分組成:決定作業在何時何地執行的編排層,持久化資料的儲存層,以及實際處理資料的計算層。
我們瞭解了分散式檔案系統與物件儲存如何透過按塊複製、快取和元資料服務來管理大檔案,以及現代批處理框架如何透過可插拔 API 與這些系統互動。我們還討論了編排器如何在大型叢集中排程任務、分配資源和處理故障,並比較了排程單個作業的作業編排器,與管理一組依賴圖作業整個生命週期的工作流編排器。
我們考察了幾種批處理模型,首先是 MapReduce 及其經典的 map 和 reduce 函式,隨後轉向 Spark、Flink 等資料流引擎;它們的資料流 API 更易使用,效能也更好。為了理解批處理作業如何擴充套件,我們還介紹了混洗(shuffle)演算法——這項基礎操作使分組、連線和聚合成為可能。
隨著批處理系統逐漸成熟,關注點也轉向了易用性。SQL 等高階查詢語言和資料框 API 降低了批處理作業的使用門檻,也讓它們更容易最佳化。查詢最佳化器會把宣告式查詢轉換成高效的執行計劃。
最後我們回顧了批處理常見用例:
- ETL 資料管道:透過定時工作流在不同系統之間提取、轉換和載入資料;
- 分析:批處理作業既支援預聚合的儀表板,也支援即席查詢;
- 機器學習:批處理作業負責準備和處理大規模訓練資料集;
- 用批處理輸出填充面向生產流量的系統:通常經由流或批次載入工具,把衍生資料提供給使用者。
下一章我們將轉向流處理,其中的輸入是 無界的(unbounded):作業依然存在,但它的輸入是永無止境的資料流。由於任何時刻都可能有更多工作到來,作業永遠不會完成。我們會看到,流處理與批處理在某些方面相似,但“流是無界的”這一假設也會在很大程度上改變系統的構建方式。
腳註
參考文獻
Nathan Marz. How to Beat the CAP Theorem. nathanmarz.com, October 2011. Archived at perma.cc/4BS9-R9A4 ↩︎
Molly Bartlett Dishman and Martin Fowler. Agile Architecture. At O’Reilly Software Architecture Conference, March 2015. ↩︎
Jeffrey Dean and Sanjay Ghemawat. MapReduce: Simplified Data Processing on Large Clusters. At 6th USENIX Symposium on Operating System Design and Implementation (OSDI), December 2004. ↩︎ ↩︎
Shivnath Babu and Herodotos Herodotou. Massively Parallel Databases and MapReduce Systems. Foundations and Trends in Databases, volume 5, issue 1, pages 1–104, November 2013. doi:10.1561/1900000036 ↩︎
David J. DeWitt and Michael Stonebraker. MapReduce: A Major Step Backwards. Originally published at databasecolumn.vertica.com, January 2008. Archived at perma.cc/U8PA-K48V ↩︎
Henry Robinson. The Elephant Was a Trojan Horse: On the Death of Map-Reduce at Google. the-paper-trail.org, June 2014. Archived at perma.cc/9FEM-X787 ↩︎
Urs Hölzle. R.I.P. MapReduce. After having served us well since 2003, today we removed the remaining internal codebase for good. twitter.com, September 2019. Archived at perma.cc/B34T-LLY7 ↩︎
Adam Drake. Command-Line Tools Can Be 235x Faster than Your Hadoop Cluster. aadrake.com, January 2014. Archived at perma.cc/87SP-ZMCY ↩︎
sort: Sort text files. GNU Coreutils 9.7 Documentation, Free Software Foundation, Inc., 2025. ↩︎Michael Ovsiannikov, Silvius Rus, Damian Reeves, Paul Sutter, Sriram Rao, and Jim Kelly. The Quantcast File System. Proceedings of the VLDB Endowment, volume 6, issue 11, pages 1092–1101, August 2013. doi:10.14778/2536222.2536234 ↩︎
Andrew Wang, Zhe Zhang, Kai Zheng, Uma Maheswara G., and Vinayakumar B. Introduction to HDFS Erasure Coding in Apache Hadoop. blog.cloudera.com, September 2015. Archived at archive.org ↩︎
Andy Warfield. Building and operating a pretty big storage system called S3. allthingsdistributed.com, July 2023. Archived at perma.cc/7LPK-TP7V ↩︎
Vinod Kumar Vavilapalli, Arun C. Murthy, Chris Douglas, Sharad Agarwal, Mahadev Konar, Robert Evans, Thomas Graves, Jason Lowe, Hitesh Shah, Siddharth Seth, Bikas Saha, Carlo Curino, Owen O’Malley, Sanjay Radia, Benjamin Reed, and Eric Baldeschwieler. Apache Hadoop YARN: Yet Another Resource Negotiator. At 4th Annual Symposium on Cloud Computing (SoCC), October 2013. doi:10.1145/2523616.2523633 ↩︎
Richard M. Karp. Reducibility Among Combinatorial Problems. Complexity of Computer Computations. The IBM Research Symposia Series. Springer, 1972. doi:10.1007/978-1-4684-2001-2_9 ↩︎
J. D. Ullman. NP-Complete Scheduling Problems. Journal of Computer and System Sciences, volume 10, issue 3, June 1975. doi:10.1016/S0022-0000(75)80008-0 ↩︎
Gilad David Maayan. The complete guide to spot instances on AWS, Azure and GCP. datacenterdynamics.com, March 2021. Archived at archive.org ↩︎
Abhishek Verma, Luis Pedrosa, Madhukar Korupolu, David Oppenheimer, Eric Tune, and John Wilkes. Large-Scale Cluster Management at Google with Borg. At 10th European Conference on Computer Systems (EuroSys), April 2015. doi:10.1145/2741948.2741964 ↩︎
Matei Zaharia, Mosharaf Chowdhury, Tathagata Das, Ankur Dave, Justin Ma, Murphy McCauley, Michael J. Franklin, Scott Shenker, and Ion Stoica. Resilient Distributed Datasets: A Fault-Tolerant Abstraction for In-Memory Cluster Computing. At 9th USENIX Symposium on Networked Systems Design and Implementation (NSDI), April 2012. ↩︎ ↩︎
Paris Carbone, Stephan Ewen, Seif Haridi, Asterios Katsifodimos, Volker Markl, and Kostas Tzoumas. Apache Flink™: Stream and Batch Processing in a Single Engine. Bulletin of the IEEE Computer Society Technical Committee on Data Engineering, volume 38, issue 4, December 2015. Archived at perma.cc/G3N3-BKX5 ↩︎ ↩︎ ↩︎
Mark Grover, Ted Malaska, Jonathan Seidman, and Gwen Shapira. Hadoop Application Architectures. O’Reilly Media, 2015. ISBN: 978-1-491-90004-8 ↩︎ ↩︎
Jules S. Damji, Brooke Wenig, Tathagata Das, and Denny Lee. Learning Spark, 2nd Edition. O’Reilly Media, 2020. ISBN: 978-1492050049 ↩︎
Michael Isard, Mihai Budiu, Yuan Yu, Andrew Birrell, and Dennis Fetterly. Dryad: Distributed Data-Parallel Programs from Sequential Building Blocks. At 2nd European Conference on Computer Systems (EuroSys), March 2007. doi:10.1145/1272996.1273005 ↩︎
Daniel Warneke and Odej Kao. Nephele: Efficient Parallel Data Processing in the Cloud. At 2nd Workshop on Many-Task Computing on Grids and Supercomputers (MTAGS), November 2009. doi:10.1145/1646468.1646476 ↩︎
Hossein Ahmadi. In-memory query execution in Google BigQuery. cloud.google.com, August 2016. Archived at perma.cc/DGG2-FL9W ↩︎ ↩︎
Tom White. Hadoop: The Definitive Guide, 4th edition. O’Reilly Media, 2015. ISBN: 978-1-491-90163-2 ↩︎ ↩︎
Fabian Hüske. Peeking into Apache Flink’s Engine Room. flink.apache.org, March 2015. Archived at perma.cc/44BW-ALJX ↩︎
Mostafa Mokhtar. Hive 0.14 Cost Based Optimizer (CBO) Technical Overview. hortonworks.com, March 2015. Archived on archive.org ↩︎
Michael Armbrust, Reynold S. Xin, Cheng Lian, Yin Huai, Davies Liu, Joseph K. Bradley, Xiangrui Meng, Tomer Kaftan, Michael J. Franklin, Ali Ghodsi, and Matei Zaharia. Spark SQL: Relational Data Processing in Spark. At ACM International Conference on Management of Data (SIGMOD), June 2015. doi:10.1145/2723372.2742797 ↩︎
Kaya Kupferschmidt. Spark vs Pandas, part 2 – Spark. towardsdatascience.com, October 2020. Archived at perma.cc/5BRK-G4N5 ↩︎
Ammar Chalifah. Tracking payments at scale. bolt.eu.com, June 2025. Archived at perma.cc/Q4KX-8K3J ↩︎
Nafi Ahmet Turgut, Hamza Akyıldız, Hasan Burak Yel, Mehmet İkbal Özmen, Mutlu Polatcan, Pinar Baki, and Esra Kayabali. Demand forecasting at Getir built with Amazon Forecast. aws.amazon.com.com, May 2023. Archived at perma.cc/H3H6-GNL7 ↩︎
Jason (Siyu) Zhu. Enhancing homepage feed relevance by harnessing the power of large corpus sparse ID embeddings. linkedin.com, August 2023. Archived at archive.org ↩︎
Avery Ching, Sital Kedia, and Shuojie Wang. Apache Spark @Scale: A 60 TB+ production use case. engineering.fb.com, August 2016. Archived at perma.cc/F7R5-YFAV ↩︎
Edward Kim. How ACH works: A developer perspective — Part 1. engineering.gusto.com, April 2014. Archived at perma.cc/F67P-VBLK ↩︎
Zhamak Dehghani. How to Move Beyond a Monolithic Data Lake to a Distributed Data Mesh. martinfowler.com, May 2019. Archived at perma.cc/LN2L-L4VC ↩︎
Chris Riccomini. What the Heck is a Data Mesh?! cnr.sh, June 2021. Archived at perma.cc/NEJ2-BAX3 ↩︎
Chad Sanderson, Mark Freeman, B. E. Schmidt. Data Contracts. O’Reilly Media, 2025. ISBN: 9781098157623 ↩︎
Daniel Abadi. Data Fabric vs. Data Mesh: What’s the Difference? starburst.io, November 2021. Archived at perma.cc/RSK3-HXDK ↩︎
Michael Armbrust, Ali Ghodsi, Reynold Xin, and Matei Zaharia. Lakehouse: A New Generation of Open Platforms that Unify Data Warehousing and Advanced Analytics. At 11th Annual Conference on Innovative Data Systems Research (CIDR), January 2021. ↩︎
Leslie G. Valiant. A Bridging Model for Parallel Computation. Communications of the ACM, volume 33, issue 8, pages 103–111, August 1990. doi:10.1145/79173.79181 ↩︎
Stephan Ewen, Kostas Tzoumas, Moritz Kaufmann, and Volker Markl. Spinning Fast Iterative Data Flows. Proceedings of the VLDB Endowment, volume 5, issue 11, pages 1268-1279, July 2012. doi:10.14778/2350229.2350245 ↩︎
Grzegorz Malewicz, Matthew H. Austern, Aart J. C. Bik, James C. Dehnert, Ilan Horn, Naty Leiser, and Grzegorz Czajkowski. Pregel: A System for Large-Scale Graph Processing. At ACM International Conference on Management of Data (SIGMOD), June 2010. doi:10.1145/1807167.1807184 ↩︎
Richard MacManus. OpenAI Chats about Scaling LLMs at Anyscale’s Ray Summit. thenewstack.io, September 2023. Archived at perma.cc/YJD6-KUXU ↩︎
Jay Kreps. Why Local State is a Fundamental Primitive in Stream Processing. oreilly.com, July 2014. Archived at perma.cc/P8HU-R5LA ↩︎
Félix GV. Open Sourcing Venice – LinkedIn’s Derived Data Platform. linkedin.com, September 2022. Archived at archive.org ↩︎