Rust 算完之後,為什麼還要把 Arrow batch 轉回 Spark 認的列式?
這件事叫 C2R(columnar-to-row),是所有 native 加速器都跑不掉成本
有趣的地方不在「有沒有這條成本」,而在 Comet 試著把這條成本從 JVM 搬到 Rust 之後,發生了什麼事。
Spark 的物理計畫是一棵tree
Comet 攔截到的通常是子樹或若干運算,很少是整棵樹。
表示同一次執行裡,有些節點在 Comet 的管線內(資料是 Arrow 列式),有些節點還是 Spark 自己跑(資料是 InternalRow / UnsafeRow)。兩者交界處,資料必須從 Arrow 轉回 Spark 行式。這就是 C2R。
真正不需要 C2R 的情況其實不多:
注意 collect() 不在這個名單裡。就算整條 pipeline 都是 native,結果要回到 driver 的那一刻還是得 row 化。所以 C2R 是常態,不是例外。
換個角度看,D25 講「Arrow C 是替換邊界的契約」的時候,其實已經隱含了這件事:只要邊界不吞掉整條 pipeline,出邊界的那一刻,格式就要轉回外面認的樣子。
C2R 的成本大致就是「要轉多少資料」除以「轉多快」
native 那段算得越快,整體時間裡剩下的就越多是 C2R,C2R 的絕對時間沒有變,但它的相對比重會被推上去,也就是說,加速器越強,回程就越顯眼。
這是一個很合理的動機:既然回程遲早會變瓶頸,那不如提早把它優化掉。Comet 就是這樣想的。
原本 Comet 的 C2R 在 JVM 側做(CometColumnarToRowExec)。後來加了一個 Rust 版本(CometNativeColumnarToRowExec),從 0.14 起改成預設啟用。
Rust 版本直接照著 Spark 的 UnsafeRow 記憶體佈局寫 bytes。這件事做得到,是因為 UnsafeRow 的格式是公開的,不需要透過 Spark 的 API 才能產生。
這裡先澄清一個容易誤會的點:native C2R 不是 Comet 的發明。Gluten 的 C2R 和 R2C 一直都是 native 實作,Velox 裡有一個 UnsafeRowFast 專門做這件事。Comet 加 fixed-width 快路徑的那個 commit,訊息第一行就寫著 Inspired by Velox UnsafeRowFast。兩條路在這件事上站在同一個位置,Comet 只是補上原本沒有的部分。
C2R 是 native 加速器的固定稅,只要替換邊界不吞掉整條 pipeline,出邊界那一刻就要退回 Spark 認的列式,而且加速器越快,這條稅的相對比重越大。
那就明天見~