iT邦幫忙

2026 iThome 鐵人賽

DAY 24
0
Software Development

1+1+1>3 ~ Spark 與 DataFusion、Comet 效能煉金術 ~系列 第 24

Day 24 Comet codegen dispatcher:getSupportLevel、canHandle 與修法時機

  • 分享至 

  • xImage
  •  

D21 講了 Comet 表達式的三條路,第三條 dispatcher 是 opt-in 的補救,今天介紹兩件 D21 沒展開的事:

  1. Dispatcher 自己也有接不住的 case,canHandle 這道閘門在哪、擋掉什麼
  2. 一個 serde 忘了開 opt-in 導致白白 fallback,修法為什麼只能寫在 getSupportLevel,不能在 convert.orElse 一個 dispatcher 回來救

第二點是全篇的核心,它揭露一個很反直覺的事實:fallback 不只是效能決策,還是相容性合約的一部分。誰在什麼時機放棄,會改變上游祖先節點算出來的答案。

一個具體場景

timestamp_seconds(x) 這個表達式,Comet 目前的 serde 只接受四種輸入型別:IntegerTypeLongTypeFloatTypeDoubleType。碰到 DecimalTypeByteTypeShortType 直接回 Unsupported

Rust 那邊 native kernel 的 signature 也只列了 Int32Int64Float32Float64,decimal 或 int8 / int16 進來就是 signature 不匹配。

問題不在「宣告不支援」,而在不支援的代價
Unsupported 又沒有其他出路時,整個 projection 退回 Spark,代價不只這一個表達式變慢,而是整段 native pipeline 被切開,中間插進 columnar 與 row 之間的 transition。這就是 D21 講過的白白 fallback

修法看起來很小:讓這個 serde 混入 CodegenDispatchFallbackUnsupported 就會走 dispatcher 而不是掉回 Spark,加一行 with CodegenDispatchFallback 就結束。

為什麼是這一行、寫在這個位置、而不是別的地方,是接下來要拆的兩件事。

Dispatcher 也不是萬能

Dispatcher 的入口是 CometScalaUDF.emitJvmCodegenDispatch,它做的事情大致是:把 Spark expression 綁定成 BoundReference 樹,用 Spark 自己的 closure serializer 序列化成 bytes,包成一顆吃 Arrow batch 的 codegen kernel。求值在 JVM,但資料維持 Arrow 格式,語意跟 Spark 逐位元組相同

聽起來 dispatcher 就是萬能兜底,但事實不是這樣。
CometBatchKernelCodegen.canHandle 是一道 plan-time 的閘門,只要它回 Some(reason),dispatcher 就拒收,乾淨地退回 Spark,而不是讓 Janino 在執行期炸掉。它擋的東西大致三類:

輸出型別必須過 isSupportedDataType,這個集合的邊界就是能跨 Arrow FFI 的型別NullTypeObjectType、Variant 都不在裡面,複合型別要遞迴檢查 children。

巢狀欄位總數不能超過 spark.sql.codegen.maxFields
這是照抄 Spark WholeStageCodegenExec 的上限,避免生出的 Janino 類別大到 JIT 或 classloader 出事。

child expression 也要能被 dispatcher 處理,這是遞迴的。

意義是什麼?
D21 問過「為什麼不 default-on 全部 serde」,當時的答案是 JNI 成本跟語意漏洞會被掩蓋。今天多一條:dispatcher 本身有 canHandle 閘門,本來就有一部分 case 它接不住。就算全部 opt-in,那些 case 還是得走 fallback。default-on 並不能真的根絕 fallback,只是把「開發者何時被迫思考這件事」從 serde 撰寫時延後到執行時。

為什麼修法是 mixin 不是 .orElse

現在到核心,既然 dispatcher 已經寫好了,看起來有兩種寫法可以讓它接手一個原本會 fallback 的 serde:

A(mixin):serde 加一行 with CodegenDispatchFallback。框架在 exprToProtoInternal 看到 getSupportLevelUnsupported,會呼叫 dispatchIfFallback,發現這個 serde 有貼紙就進 dispatcher。

B(.orElse in convert)convertNone.orElse(dispatcher()),一行 fallback。

B 看起來更集中、更直接,但 B 是錯的,原因不是程式碼風格,是放棄的時機

getSupportLevelUnsupported 發生在任何 ancestor 轉換之前。這時候放棄,整棵子樹連同祖先一起乾淨地退回 Spark,祖先根本沒開始 native 轉換。

convertNone 就完全不同了。此時祖先節點已經在「這個子節點會 fallback,我要走 Spark 相容分支」的前提下開始 native 轉換。如果 dispatcher 在這個時機突然接手把子節點保住,祖先原本依賴的前提被推翻了:它以為自己要處理一個 Spark 產出的結果,實際上收到的是 dispatcher 產出的結果,兩者格式相同但祖先的行為分支不對。答案就會錯。

換句話說:

  • A 在祖先做決定之前放棄,祖先看到的世界是一致的
  • B 在祖先已經做完決定之後才反悔,祖先看到的世界跟事實不符

Fallback 這件事是相容性合約的一部分,不只是效能取捨。誰在什麼時機放棄,會改變上游節點算出來的答案。

這條規則寫成一句就是:要救 fallback,只能在 getSupportLevel 這一格救
convert 裡的補救永遠是 unsound 的

三張型別支援表要對齊

想通了上一段,還有一個現實問題:加了 mixin 就一定能救嗎?不一定

一個表達式在 Comet 有三份型別清單,分別長在不同地方:

清單 位置 用途
serde getSupportLevel serde object 這個表達式接受什麼輸入型別
native Signature native/spark-expr/ Rust kernel 認得什麼型別
dispatcher isSupportedDataType CometBatchKernelCodegen Arrow FFI 能跨什麼型別

一張 fallback 案例能靠「加一行 mixin」修好的前提是:第三份比前兩份寬。也就是 native 不認得的型別,dispatcher 認得。

反過來如果 dispatcher 也不認得(例如 Variant,或欄位數超過 maxFields),加了 mixin 也沒用。這時要嘛動 native kernel 讓它接更多型別,要嘛在 serde 端補一個 cast 把型別降到 native 認得的範圍。

判斷「這個 fallback 是 mixin 就能救、還是要動 native」,就是看這三張表能不能對齊。

結論

Dispatcher 是中間路線,但它有自己接不住的 case(canHandle 閘門),也有自己該寫的位置(getSupportLevel 這一格)。修法選 mixin 不選 .orElse,是因為放棄的時機決定祖先節點的正確性。這條合約寫在 serde,不能拖到 convert。要判斷一個 fallback 案例能不能用 mixin 修好,就看 serde / native signature / dispatcher isSupportedDataType 這三張表能不能對齊

那就明天見~

參考資料


上一篇
Day 23 Comet protobuf vs Substrait:自訂 IR 還是跨引擎標準
下一篇
Day 25 Arrow C Data Interface 與 JNI 邊界
系列文
1+1+1>3 ~ Spark 與 DataFusion、Comet 效能煉金術 ~26
圖片
  熱門推薦
圖片
{{ item.channelVendor }} | {{ item.webinarstarted }} |
{{ formatDate(item.duration) }}
直播中

尚未有邦友留言

立即登入留言