D21 講了 Comet 表達式的三條路,第三條 dispatcher 是 opt-in 的補救,今天介紹兩件 D21 沒展開的事:
canHandle 這道閘門在哪、擋掉什麼getSupportLevel,不能在 convert 裡 .orElse 一個 dispatcher 回來救第二點是全篇的核心,它揭露一個很反直覺的事實:fallback 不只是效能決策,還是相容性合約的一部分。誰在什麼時機放棄,會改變上游祖先節點算出來的答案。
timestamp_seconds(x) 這個表達式,Comet 目前的 serde 只接受四種輸入型別:IntegerType、LongType、FloatType、DoubleType。碰到 DecimalType、ByteType、ShortType 直接回 Unsupported。
Rust 那邊 native kernel 的 signature 也只列了 Int32、Int64、Float32、Float64,decimal 或 int8 / int16 進來就是 signature 不匹配。
問題不在「宣告不支援」,而在不支援的代價。
回 Unsupported 又沒有其他出路時,整個 projection 退回 Spark,代價不只這一個表達式變慢,而是整段 native pipeline 被切開,中間插進 columnar 與 row 之間的 transition。這就是 D21 講過的白白 fallback。
修法看起來很小:讓這個 serde 混入 CodegenDispatchFallback,Unsupported 就會走 dispatcher 而不是掉回 Spark,加一行 with CodegenDispatchFallback 就結束。
為什麼是這一行、寫在這個位置、而不是別的地方,是接下來要拆的兩件事。
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 的型別。NullType、ObjectType、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 撰寫時延後到執行時。
.orElse現在到核心,既然 dispatcher 已經寫好了,看起來有兩種寫法可以讓它接手一個原本會 fallback 的 serde:
A(mixin):serde 加一行 with CodegenDispatchFallback。框架在 exprToProtoInternal 看到 getSupportLevel 回 Unsupported,會呼叫 dispatchIfFallback,發現這個 serde 有貼紙就進 dispatcher。
B(.orElse in convert):convert 回 None 時 .orElse(dispatcher()),一行 fallback。
B 看起來更集中、更直接,但 B 是錯的,原因不是程式碼風格,是放棄的時機。
getSupportLevel 回 Unsupported 發生在任何 ancestor 轉換之前。這時候放棄,整棵子樹連同祖先一起乾淨地退回 Spark,祖先根本沒開始 native 轉換。
convert 回 None 就完全不同了。此時祖先節點已經在「這個子節點會 fallback,我要走 Spark 相容分支」的前提下開始 native 轉換。如果 dispatcher 在這個時機突然接手把子節點保住,祖先原本依賴的前提被推翻了:它以為自己要處理一個 Spark 產出的結果,實際上收到的是 dispatcher 產出的結果,兩者格式相同但祖先的行為分支不對。答案就會錯。
換句話說:
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 這三張表能不能對齊
那就明天見~
serde/CometScalaUDF.scala(emitJvmCodegenDispatch 入口): https://github.com/apache/datafusion-comet/blob/main/spark/src/main/scala/org/apache/comet/serde/CometScalaUDF.scala
serde/CometExpressionSerde.scala(CodegenDispatchFallback trait): https://github.com/apache/datafusion-comet/blob/main/spark/src/main/scala/org/apache/comet/serde/CometExpressionSerde.scala
serde/QueryPlanSerde.scala(exprToProtoInternal 分流、dispatchIfFallback): https://github.com/apache/datafusion-comet/blob/main/spark/src/main/scala/org/apache/comet/serde/QueryPlanSerde.scala
serde/datetime.scala(具體場景 CometSecondsToTimestamp): https://github.com/apache/datafusion-comet/blob/main/spark/src/main/scala/org/apache/comet/serde/datetime.scala
WholeStageCodegenExec(spark.sql.codegen.maxFields 出處): https://github.com/apache/spark/blob/master/sql/core/src/main/scala/org/apache/spark/sql/execution/WholeStageCodegenExec.scala
isSupportedDataType 的邊界依據): https://arrow.apache.org/docs/format/CDataInterface.html