昨天講單一表達式的三條路,其中 dispatcher 是為了避免「一個表達式沒 native、整個 Project 白白 fallback」。今天上一層到 operator:一整個 operator 沒轉成功時,Comet 怎麼記錄理由?
Comet 用 Spark 的 extension 機制掛上自己的規則(CometSparkSessionExtensions.scala),Spark 走到 physical plan 之後會呼叫這些規則,遍歷 plan 樹問每個節點「能不能換成 Comet 的版本」。能就換,不能就原封不動留給 Spark。
留給 Spark 這件事就叫 fallback。這裡有個反直覺的認知要先建立:fallback 是正常且必要的設計。Comet 不可能支援 Spark 全部功能,也不會想支援。使用者需要看到的不是「有沒有 fallback」,是「哪些節點在 fallback、為什麼」。
每一種 Spark operator 對應一個 serde 物件,註冊在 CometExecRule.scala 的兩張表:nativeExecs 是完全原生執行的 operator,sinks 是外層仍在 JVM、內部包一個原生 ScanExec 的 operator。差別在資料是否走完全部 native 管線,D26 會再細講。
每個 serde 有四個方法:
| 方法 | 負責什麼 |
|---|---|
getClassName(op) |
EXPLAIN 該印什麼名字 |
getSupportLevel(op) |
這個節點能不能轉?不能的話原因是什麼 |
convert(op, builder, childOp*) |
產生 protobuf |
createExec(nativeOp, op) |
包成 Comet 的 SparkPlan 節點 |
SupportLevel 是一個 sealed trait,只有三個 case:
| 值 | 意思 | 預設行為 |
|---|---|---|
Compatible |
與 Spark 結果完全一致 | 用 Comet 跑 |
Incompatible |
能跑,但結果可能與 Spark 不同 | 預設 fallback,要設 allowIncompatible 才會用 |
Unsupported |
不能跑 | fallback |
這是 Comet 最重要的契約:寧可慢,不可錯。Incompatible 預設不啟用,是使用者信任 Comet 加速結果的基礎。跟 D21 的 dispatcher 是同一套三態,只是 dispatcher 問「這條走不了時有沒有第二條」,SupportLevel 問「這條走不走得了」。
精妙的設計在這裡。QueryPlanSerde 呼叫的地方大致是:
handler.getSupportLevel(op) match {
case Unsupported(notes) =>
withFallbackReason(op, notes.getOrElse("")) // 框架自動記錄理由
false
case Incompatible(notes) if !allowIncompatible =>
withFallbackReason(op, notes.getOrElse(""))
false
case Compatible(_, _) | Incompatible(_) =>
true
}
只要 serde 回一個 Unsupported(Some(理由)),理由自動被框架掛到 plan 節點上。反觀 convert:serde 可以在裡面做任意檢查然後 return None,框架不知道你為什麼失敗,得自己呼叫 withFallbackReason。忘了就是沒有解釋的 fallback,EXPLAIN 只印一句「這個節點 fallback 了」。
| 檢查寫在 | 理由誰記錄 | 後果 |
|---|---|---|
getSupportLevel |
框架自動 | 一定有解釋 |
convert |
呼叫者自己 | 忘了就沒有解釋 |
推論:任何「光看節點本身就能決定」的檢查都該放在 getSupportLevel。convert 只該處理需要子節點轉換結果才知道的失敗。
第一層,TreeNodeTag 存理由。 Comet 用 Spark TreeNode 的 tag 機制存 fallback 理由(ExtendedExplainInfo.scala):
val FALLBACK_REASONS = TreeNodeTag[Set[String]]("FALLBACK_REASONS")
用集合而非單一字串,因為同一個節點可能有多個理由。
第二層,表達式理由上捲。 表達式樹跟 plan 樹是分開的,EXPLAIN 只走 plan 節點。理由掛在表達式上讀者永遠看不到,所以攔截結束後會把表達式的理由上捲到最近的 plan operator。這解釋了一個怪現象:EXPLAIN 印在 Project 旁邊的理由,可能其實是底下某個表達式的問題。
第三層,reportUnexplainedFallback。 前兩層都做了還是可能有節點沒理由:
if (op.children.forall(_.isInstanceOf[CometNativeExec]) && !hasFallbackReason(op)) {
if (COMET_STRICT_FALLBACK_REASONS.get(op.conf)) {
throw new IllegalStateException("... recorded no fallback reason ...")
}
withFallbackReason(op, s"${op.nodeName} is not supported")
}
條件是「所有子節點都已經是 Comet native、只有這個節點掉出去、又沒理由」。這是最可疑的一種 fallback,代表本來有機會延續 native 管線卻莫名中斷。處理方式看 spark.comet.explain.fallback.strict.enabled:production 預設 false,掛上通用訊息;Comet 自己的測試預設 true,直接拋例外。用意是讓「忘了寫理由」在開發階段就爆炸,而不是出貨給使用者一句廢話。
ObjectHashAggregate 在 Spark 裡會被拆成兩階段:
ObjectHashAggregate(Final) 每個 key 的最終結果
+- ShuffleExchange 依 group key 重新分區
+- ObjectHashAggregate(Partial) 每個分區先局部聚合
+- Scan
流過 shuffle 的不是最終值,是聚合緩衝區(例如 avg 傳的是 sum 與 count)。Comet 與 Spark 的緩衝區格式不一定相容,若 partial 用 Comet 跑而 final 用 Spark 跑,Spark 拿到看不懂的緩衝區會出錯。
所以「shuffle 未啟用時要 fallback」不是效能取捨,是正確性守衛。理由訊息要寫成不分 stage,因為 partial 跟 final 都會被同一個守衛擋下。這個檢查光看節點自己加上 config 就能決定,天生就該待在 getSupportLevel。
Comet 的 fallback 不是失敗處理,是設計決策,而診斷子系統的存在是為了讓 fallback 可被觀測。同一個守衛寫在 convert 跟寫在 getSupportLevel 只差幾行位置,但一個是沉默的 fallback,一個是有解釋的 fallback。
明天 D23 換個角度:operator 怎麼被翻成跨語言的表示塞到 Rust 那邊?Comet 為什麼不走 Substrait 這條看起來更標準的路,而選擇自訂 protobuf?
那就明天見~
CometSparkSessionExtensions.scala: https://github.com/apache/datafusion-comet/blob/main/spark/src/main/scala/org/apache/comet/CometSparkSessionExtensions.scala
CometExecRule.scala: https://github.com/apache/datafusion-comet/blob/main/spark/src/main/scala/org/apache/comet/rules/CometExecRule.scala
SupportLevel.scala: https://github.com/apache/datafusion-comet/blob/main/spark/src/main/scala/org/apache/comet/serde/SupportLevel.scala
QueryPlanSerde.scala: https://github.com/apache/datafusion-comet/blob/main/spark/src/main/scala/org/apache/comet/serde/QueryPlanSerde.scala
ExtendedExplainInfo.scala: https://github.com/apache/datafusion-comet/blob/main/spark/src/main/scala/org/apache/comet/ExtendedExplainInfo.scala
CometConf.scala: https://github.com/apache/datafusion-comet/blob/main/common/src/main/scala/org/apache/comet/CometConf.scala
interfaces.scala(AggregateMode): https://github.com/apache/spark/blob/master/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/expressions/aggregate/interfaces.scala