有讀者提問:Physical Planner 建立資料來源運算子時,Schema 到底扮演什麼角色?
這裡可以先釐清一點:Schema 並不是直接被轉換成資料來源運算子,而是建立運算子時所依據的欄位規格。
以 orders 表為例,TableProvider 會提供:
city: Utf8
amount: Float64
is_member: Boolean
DataFusion 在建立 Logical Plan 時,會利用這份 Schema 確認欄位是否存在、資料型別是否正確,以及查詢需要哪些欄位。
接著,Physical Planner 會根據 TableScan、Schema 與 TableProvider 提供的掃描方式,建立對應的資料來源執行算子。執行時,這些算子會依照 Schema 解讀 CSV、JSON 或 Parquet 中的資料,並產生符合 Schema 的 Arrow RecordBatch:
Schema
city: Utf8
amount: Float64
is_member: Boolean
↓
Arrow RecordBatch
city: [Taipei, Taipei]
amount: [1200.0, 2300.0]
is_member: [true, true]
可以將三者的關係簡化成:
Schema
→ 描述欄位名稱與資料型別
TableProvider
→ 提供 Schema 與資料來源的掃描方式
Physical Planner
→ 根據這些資訊建立執行算子
Data Source Execution
→ 讀取資料並產生 RecordBatch
因此,Schema 同時參與兩個階段:
接下來就用一段實際 SQL,完整追蹤一段 SQL 從輸入字串到查詢結果的過程!
假設有一張 orders 表:
| city | amount | is_member |
|---|---|---|
| Taipei | 1200.0 | true |
| Taichung | 850.0 | false |
| Taipei | 2300.0 | true |
| Kaohsiung | 560.5 | false |
| Taoyuan | 1750.0 | true |
我們想找出會員訂單中,各城市的訂單總金額:
SELECT city,
SUM(amount) AS total_amount
FROM orders
WHERE is_member = true
GROUP BY city;
預期結果:
| city | total_amount |
|---|---|
| Taipei | 3500.0 |
| Taoyuan | 1750.0 |
雖然這段 SQL 很短,但 Query Engine 背後會經過多個階段。
SQL String
↓
SQL Parser
↓
AST
↓
Analyzer / Planner
↓
Logical Plan
↓
Logical Optimizer
↓
Optimized Logical Plan
↓
Physical Planner
↓
Physical Plan
↓
Execution
↓
Arrow RecordBatch
↓
Query Result

圖 1:DataFusion SQL Lifecycle Part 1——SQL 從字串輸入開始,經過 Parser、Schema / Catalog 驗證、Logical Plan 建立與最佳化,產生 Optimized Logical Plan。
本圖聚焦於 SQL 解析至 Optimized Logical Plan 的流程;Physical Plan、Execution 與 Query Result 將在後續階段介紹。
最開始,DataFusion 收到的只是一段文字:
SELECT city,
SUM(amount) AS total_amount
FROM orders
WHERE is_member = true
GROUP BY city;
此時 DataFusion 還不知道:
orders 是否存在city、amount 與 is_member 是否存在amount 的資料型別是什麼SUM(amount) 是否為合法操作is_member = true 是否為合法比較這些資訊需要在後續階段分析。
SQL Parser 會進行詞法分析與語法分析,將 SQL 文字轉換成 AST(Abstract Syntax Tree,抽象語法樹)。
可以簡化成:
SELECT
├── city
└── SUM(amount) AS total_amount
FROM orders
WHERE is_member = true
GROUP BY city
Parser 主要確認 SQL 是否符合文法,例如:
SELECT 後面是否有欄位FROM 後面是否有資料表WHERE 條件是否寫法正確GROUP BY 是否放在適當位置Parser 只負責理解語法,通常還不會確認資料表或欄位是否真的存在。
例如:
SELECT unknown_column
FROM orders;
這段 SQL 可能符合語法,但 unknown_column 是否存在,要等到後續分析階段才會確認。
接著,DataFusion 會透過 Catalog 與 TableProvider 取得 orders 的 Schema:
orders
├── city: Utf8
├── amount: Float64
└── is_member: Boolean
Analyzer 會將 SQL 中的欄位名稱解析成實際欄位:
city → orders.city
amount → orders.amount
is_member → orders.is_member
同時檢查:
orders 是否存在is_member = true 是否為合法比較SUM(amount) 是否支援 Float64
GROUP BY city 是否符合聚合規則如果欄位不存在或型別不相容,查詢會在此階段失敗。
分析完成後,DataFusion 會建立 Logical Plan:
Projection
Aggregate
Filter
TableScan
展開來看:
Projection:
city
SUM(amount) AS total_amount
↓
Aggregate:
GROUP BY city
SUM(amount)
↓
Filter:
is_member = true
↓
TableScan:
orders
Logical Plan 描述的是:
這段查詢需要完成哪些工作。
此時還沒有決定要使用哪個具體執行元件。
DataFusion 會套用 Logical Optimizer Rules,嘗試改善查詢計畫。
原始查詢只需要:
city
amount
is_member
因此其他不需要的欄位可以不讀取。
WHERE is_member = true 可以盡早執行:
先篩選資料
再進行 GROUP BY
這樣可以讓較少的資料進入後續聚合階段。
最佳化後,計畫仍然描述查詢邏輯,但已經減少不必要的欄位與資料處理。
Physical Planner 會將 Optimized Logical Plan 轉換成 Physical Plan。
概念上可能變成:
ProjectionExec
AggregateExec: mode=Final
RepartitionExec: Hash(city)
AggregateExec: mode=Partial
FilterExec: is_member = true
DataSourceExec
這些節點代表實際執行時使用的元件:
| Logical Operator | Physical Operator |
|---|---|
TableScan |
DataSourceExec 或 MemoryExec |
Filter |
FilterExec |
Aggregate |
AggregateExec |
Projection |
ProjectionExec |
實際名稱會依資料來源與 DataFusion 版本而有所不同。
執行開始後,資料來源執行算子會讀取實際資料。
如果資料來自:
CSV
JSON
Parquet
資料來源元件會負責:
RecordBatch
例如讀入的資料可能表示成:
city: [Taipei, Taichung, Taipei]
amount: [1200.0, 850.0, 2300.0]
is_member: [true, false, true]
FilterExec 執行:
WHERE is_member = true
資料會被篩選成:
city: [Taipei, Taipei]
amount: [1200.0, 2300.0]
is_member: [true, true]
Taichung 的資料因為 is_member = false,不會繼續傳遞到下游。
接著執行:
GROUP BY city
以及:
SUM(amount)
如果資料被切成多個 Partition,DataFusion 可能將聚合分成兩個階段:
Partial Aggregate
↓
Repartition
↓
Final Aggregate
在這個範例中:
Taipei → 1200.0 + 2300.0 = 3500.0
Taoyuan → 1750.0
Partial Aggregate 先在各分區計算部分結果,Final Aggregate 再合併這些結果。
最後,ProjectionExec 只保留查詢要求的欄位:
city
total_amount
執行結果可以表示成 Arrow RecordBatch:
city: [Taipei, Taoyuan]
total_amount: [3500.0, 1750.0]
CLI 或應用程式再將這些資料格式化成表格,顯示給使用者。
在 DataFusion CLI 中執行:
EXPLAIN FORMAT INDENT
SELECT city,
SUM(amount) AS total_amount
FROM orders
WHERE is_member = true
GROUP BY city;
可能會看到類似的結果:
logical_plan
Projection
Aggregate
Filter
TableScan
physical_plan
ProjectionExec
AggregateExec
RepartitionExec
AggregateExec
FilterExec
DataSourceExec
實際輸出會受到以下因素影響:
因此,上面的計畫是理解流程的示意,不代表所有環境都會完全相同。
現在可以回答前幾天讀者提出的問題:
Planning 和 Execution 是獨立的兩塊嗎?
它們是不同階段,但不是互相獨立。
Planning
負責產生「如何執行」的計畫
↓
Execution
依照這份計畫實際處理資料
Planning 階段如果沒有產生 Physical Plan,Execution Engine 就不知道要使用哪些算子與執行順序。
而 Execution 階段則會將 Physical Plan 轉化成實際的資料流處理。
今天完整追蹤了一段 SQL 的生命週期:
RecordBatch 在算子之間傳遞可以用一句話總結:
Logical Plan 說明查詢要完成什麼,Physical Plan 決定查詢如何執行,而 Execution Engine 依照這份計畫處理實際資料。
接下來將進入第四階段,介紹 Query Optimizer 如何透過欄位裁剪、條件下推與資料來源最佳化,降低查詢成本。
iThome鐵人賽