D25 講清楚了資料怎麼過界:Arrow C Data Interface 定義結構體佈局,release callback 定義所有權。今天講一個更麻煩的東西,記憶體額度怎麼過界。
Spark 的記憶體管理從 Project Tungsten 開始就有一套完整的機制:
UnifiedMemoryManager 管 executor 層級的額度,execution 與 storage 動態互相借用TaskMemoryManager 管單一 task 的額度,底下掛著一堆 MemoryConsumer
TaskMemoryManager 會去叫其他 consumer spill,把額度讓出來這是一套協作式的 spill 機制,Spark 能撐住大資料靠的就是它。
DataFusion 那邊也有一套:MemoryPool,operator 透過 MemoryReservation 跟它要額度,要不到就自己走 spill 路徑。
問題來了。Comet 把一部分執行搬到 Rust,於是同一個 executor process 裡出現兩本帳:
executor process
├── JVM 側
│ └── TaskMemoryManager
│ ├── UnsafeExternalSorter (知道自己用了多少)
│ ├── ShuffleExternalSorter (知道自己用了多少)
│ └── ???
└── Rust 側
└── DataFusion MemoryPool
├── SortExec (JVM 完全看不到)
└── AggregateExec (JVM 完全看不到)
如果兩本帳各記各的,結果是可預期的:兩邊各自覺得自己還有額度,加起來超過容器上限,然後被 OOM killer 砍掉。而且死的是整個 executor,不是一個 task。
Comet 的做法是實作一個自訂的 MemoryPool,它不自己管額度,而是每次 grow 的時候打一通 JNI 回 JVM,向 Spark 的 TaskMemoryManager 申請。這就是 unified pool 這個名字的由來。
一次 grow 的完整流程:
Rust operator 需要 8 MB
│
▼
Comet MemoryPool.grow(8 MB)
│
▼ JNI
TaskMemoryManager.acquireExecutionMemory(8 MB)
│
├─ 額度夠 ────────────────► 回傳 8 MB
│
├─ 額度不夠 ──► 叫其他 consumer spill
│ (UnsafeExternalSorter 寫 run 到磁碟)
│ └─► 回傳 8 MB
│
└─ 還是不夠 ──► 回傳 3 MB(部分核准)或 0
│
▼
Rust 側自己走 spill 路徑
這裡有三件事值得注意。
第一,這通 JNI 呼叫是會阻塞的,而且可能塞很久。 因為 JVM 那邊為了滿足你,可能正在逼其他 consumer 把資料寫到磁碟。Rust 這邊的 grow 看起來像一個普通函式呼叫,實際上背後可能是一次磁碟 I/O。
第二,回傳值可能是部分核准。 Spark 的 acquireExecutionMemory 本來就允許「我只給你這麼多」,所以 Rust 側的 pool 要處理 grow 成功但拿到的比要的少這件事,不能當成二元的成功或失敗。
第三,這條路把兩個 runtime 的執行緒綁在一起了。 Comet 的 native 執行跑在 tokio runtime 上,也就是說打這通 JNI 的是 tokio 的 worker thread。這件事後來咬了一口:jni-rs 會用 AttachCurrentThread 把這些執行緒惰性掛進 JVM,而那樣掛進去的是非 daemon 執行緒,DestroyJavaVM 會等所有非 daemon 執行緒結束才肯走。結果就是應用程式跑完之後 JVM 掛在那裡退不掉。後來的修法是改用 AttachCurrentThreadAsDaemon,跟 Spark 自己的執行緒池同樣的理由。
D12 講 tokio 的時候我們只把它當成一個 CPU 工作的排程器。現在它多了一個身分:它是會回頭呼叫 JVM 的那一方。
Comet 的 pool 是同一個 Spark task 內所有 native 執行環境共用的。
為什麼要共用?因為一個 task 裡可能同時存在不只一個 native 執行環境。最典型的是 shuffle:Comet 會同時跑兩個,一個負責 shuffle 之前的運算子,一個負責 shuffle writer。如果兩個各自拿一份完整額度,那 per-task 的上限等於形同虛設。共用一個 pool,才守得住那條線。
這也解釋了為什麼 pool 的分配策略會有兩種選擇(spark.comet.exec.memoryPool):
| 策略 | 行為 | 適合 |
|---|---|---|
fair_unified(預設) |
任何一個 reservation 不能超過 pool_size / 目前的 reservation 數 |
事先知道會有多個運算子都要 spill |
greedy_unified |
先來先得 | 不需要 spill,或只有單一個會 spill 的運算子 |
fair_unified 有一個反直覺的副作用:它有時候會在記憶體其實還夠的情況下就叫你去 spill,因為它要幫還沒出現的其他 reservation 留位置。這就是 D17 講 Greedy vs Fair 那組取捨,換到跨 runtime 的場景再演一次。
到這裡聽起來很完整,但這一層其實是 Comet 目前最不成熟的部分之一。官方 roadmap 自己把「memory accounting、reservation strategies、spill integration」列為待改進項目,理由寫得很直白:這些會減少 OOM。
具體有幾個對不準的地方。
一、計帳本身就有誤差。 tuning guide 明說 Comet 的記憶體計帳不是 100% 準確,可能用得比它登記的多,所以提供了 spark.comet.exec.memoryPool.fraction,讓你設一個小於 1.0 的係數,故意不把額度用滿。一個配置選項的存在就是承認:這本帳記得不準,請你自己留 buffer。
二、spill 的數字回報不完整。 native 側 spill 的位元組數要回報到 Spark 的 task metrics(diskBytesSpilled / memoryBytesSpilled),目前只有 native shuffle write 那條路有接上。一個會 spill 的 native 運算子如果跑在非 shuffle 的 stage,Spark UI 的 Stages 頁面會顯示這個 task 完全沒有 spill,實際上它寫了一堆到磁碟。SQL 頁面看得到,task 頁面看不到。
三、不是每個 native 運算子都會 spill。 Comet 的 native hash join 目前要求 build side 整個放得進記憶體,沒有 spill 路徑。這條也還在 roadmap 上。
這三件事加起來會長成什麼樣子?社群裡有一則 300 GB 規模的回報講得很清楚:同樣的 left join 跟 distinct,vanilla Spark 會 spill 到磁碟然後慢慢跑完,Comet 直接在 native sort 跟 native aggregation 那裡 OOM。回報者把 tuning guide 上所有記憶體選項都試過了,結果不變。
回到今天的問題。
unified pool 統一的是額度的來源。Rust 側不再自己決定能用多少,它每一次要記憶體都回 JVM 問,所以兩本帳變成一本,容器的上限守得住。這件事是成立的,而且是必要的。
但它沒有統一spill 的能力。Spark 的協作式 spill 之所以強,是因為底下每一個 consumer 都會 spill,TaskMemoryManager 要不到記憶體時總有人可以讓。Comet 這邊,有些 native 運算子會 spill、有些不會,回報的數字還不完整,計帳本身還有誤差。
共用同一本帳,不代表共用同一張安全網。 這是我覺得這一層最值得記住的一句話,也是為什麼 Comet 在大資料量下的失敗模式跟 vanilla Spark 不一樣:Spark 是變慢,Comet 是掛掉。
Spark 跟 DataFusion 各有一套記憶體管理,放在同一個 executor 裡會各記各的帳,加起來超過容器上限。Comet 的解法是讓 Rust 側的 MemoryPool 每次 grow 都透過 JNI 回頭跟 Spark 的 TaskMemoryManager 申請,額度的來源因此統一。這通呼叫會阻塞、會被部分核准,而且是從 tokio 的 worker thread 打出去的。
pool 綁在 task 上而不是 operator 上,因為一個 task 裡可能同時跑多個 native 執行環境,shuffle 就是。fair_unified 與 greedy_unified 是 D17 那組取捨的跨 runtime 版本。
但額度統一不等於安全網統一。計帳有誤差、spill 數字回報不完整、有些運算子根本不能 spill,三件事疊起來的結果是 Comet 在記憶體壓力下的失敗模式比 Spark 更硬。這一層官方自己還掛在 roadmap 上。
那就明天見~