昨天有讀者問:
DataFusion 是如何利用 Schema 產生 RecordBatch?DataFusion 和 Schema 的關係到底是什麼?
簡短的答案是:
Schema 不會憑空產生 RecordBatch。
Schema 定義資料規格;
Data Source 讀取實際資料;
DataFusion 依照 Schema 規劃與驗證查詢,
並在執行時讓資料來源分批產生符合規格的 RecordBatch。
也就是說,Schema 像資料的設計圖,RecordBatch 則是依照設計圖裝好的實際資料。
今天就從 CSV、JSON、Parquet 等資料來源出發,追蹤它們如何進入 DataFusion,最後成為 Query Engine 能處理的 Arrow RecordBatch。
假設有一份訂單資料:
| order_id | city | amount | is_member |
|---|---|---|---|
| 1 | 台北 | 1200.0 | true |
| 2 | 台中 | 850.0 | false |
| 3 | 台北 | 2300.0 | true |
資料原本可能放在 CSV、JSON 或 Parquet 檔案中。
但進入 DataFusion 後,它會以 Arrow 的欄式資料批次表示:
RecordBatch
├── order_id : [1, 2, 3]
├── city : [台北, 台中, 台北]
├── amount : [1200.0, 850.0, 2300.0]
└── is_member : [true, false, true]
DataFusion 的 Filter、Projection、Aggregate、Join 等執行節點,會接收 RecordBatch、處理資料,並輸出新的 RecordBatch。

圖 1:DataFusion 從資料來源讀取資料、產生 Arrow RecordBatch,並執行查詢的簡化流程。圖由本文整理。
這張圖刻意分成兩個階段:
規劃階段
→ 先確認資料結構,決定 SQL 要如何執行。
執行階段
→ 實際讀取資料,分批產生 RecordBatch 並完成運算。
這個區分很重要。DataFusion 不會一收到 SQL 就直接把所有檔案讀進記憶體。
假設我們執行:
SELECT city, SUM(amount) AS total_amount
FROM orders
WHERE is_member = true
GROUP BY city;
DataFusion 首先需要知道 orders 是什麼,以及它有哪些欄位。
這時候 TableProvider 會提供可查詢資料表的資訊,例如:
Table:orders
Schema:
order_id : Int64
city : Utf8
amount : Float64
is_member : Boolean
有了 Schema,DataFusion 才能驗證:
city 是否存在?
amount 是否為數值型別,能否執行 SUM?
is_member 是否可以和 true 比較?
接著,它會建立 Logical Plan:
Aggregate: GROUP BY city、SUM(amount)
↓
Filter: is_member = true
↓
Table Scan: orders
Logical Plan 描述的是「要做什麼」,而不是「現在就讀哪些資料」。
不同資料來源取得 Schema 的方式不太一樣。
| 資料格式 | Schema 的來源 |
|---|---|
| CSV | 由標頭與部分資料內容推論,或由使用者手動指定 |
| JSON / NDJSON | 由 JSON 記錄推論,或由使用者手動指定 |
| Parquet | 從檔案本身的 Metadata 讀取 |
| 自訂資料來源 | 由開發者實作並提供 |
例如 CSV 的一行資料是:
1,台北,1200.0,true
DataFusion 需要知道這些文字分別代表:
1 → Int64
台北 → Utf8
1200.0 → Float64
true → Boolean
若 CSV 推論到的型別不符合預期,或資料格式本來就不穩定,便可以手動指定 Schema。
Parquet 則已在檔案 Metadata 中保存欄位與型別資訊,所以通常不需要額外提供 Schema。
可以把 TableProvider 想成 DataFusion 與資料來源之間的橋樑。
它負責告訴 DataFusion:
這張表的 Schema 是什麼?
資料在哪裡?
能否裁剪欄位?
能否下推篩選條件?
實際掃描資料時,應該建立什麼執行計畫?
因此,SQL 裡的:
FROM orders
不是單純指向一個 orders.parquet 檔案,而是指向一個可查詢的資料表來源。
DataFusion 內建的 CSV、JSON、Parquet 讀取支援,會將檔案來源包裝成可查詢的 TableProvider。
Logical Plan 完成後,DataFusion 會建立 Physical Plan。
它決定更接近實作層面的事情,例如:
要讀取哪些檔案?
要讀取哪些欄位?
能否跳過部分資料?
資料要如何分區與平行處理?
以上面的 SQL 為例,order_id 完全沒有被使用。
因此,Physical Plan 可以只讀取:
city
amount
is_member
這就是 Day 7 提過的 Column Pruning。
若資料來源是 Parquet,DataFusion 還可能依條件與檔案統計資訊,跳過不可能符合條件的資料區塊。
這個階段仍是在規劃「如何讀取」;真正的資料 I/O 會在執行階段發生。
開始執行查詢後,DataFusion 的掃描執行節點會真正開啟檔案、讀取資料。
對內建的檔案資料來源而言,這類掃描工作可由 DataSourceExec 負責。
概念流程如下:
CSV / JSON / Parquet
↓
DataSourceExec 依 Physical Plan 讀取
↓
Arrow RecordBatch 1
↓
Arrow RecordBatch 2
↓
Arrow RecordBatch 3
↓
...
每個 RecordBatch 都必須符合預期的輸出 Schema。
例如,因為 order_id 已被裁剪,讀取到的一批資料可能是:
RecordBatch
├── city : [台北, 台中, 台北]
├── amount : [1200.0, 850.0, 2300.0]
└── is_member : [true, false, true]
這裡的重點是:
Schema 定義這批資料應有的欄位與型別。
Data Source 提供實際值。
DataSourceExec 將它們分批輸出成 RecordBatch。
接著,這批資料會流過不同的執行節點。
RecordBatch
↓
Filter:is_member = true
↓
Projection:保留 city、amount
↓
Aggregate:依 city 計算 SUM(amount)
↓
Query Result
以這批資料為例,Filter 後會得到:
city : [台北, 台北]
amount : [1200.0, 2300.0]
Aggregate 再計算:
city | total_amount
-----|-------------
台北 | 3500.0
每個執行節點都可能產生新的 RecordBatch,而且它們的 Schema 也可能改變。
例如:
掃描 orders 的 Schema
order_id, city, amount, is_member
↓
Filter 後
欄位不變,資料列變少
↓
Projection 後
city, amount
↓
Aggregate 後
city, total_amount
所以 Schema 不只是讀檔時使用一次;它會伴隨整個 Query Plan,協助 DataFusion 知道每個節點應該接收與輸出什麼資料。
三種格式的原始資料差異很大:
CSV
→ 純文字,通常需要推論或指定型別
JSON
→ 可有巢狀資料、缺失欄位或結構變化
Parquet
→ 原生欄式格式,帶有 Schema 與 Metadata
但它們被讀進 DataFusion 後,都會轉成 Arrow RecordBatch。
CSV ─┐
JSON ─┼→ Arrow RecordBatch → DataFusion Execution Operators
Parquet ─┘
後續的 Filter、Projection、Aggregate、Join 不需要為 CSV、JSON、Parquet 各自實作一次。
這正是 Arrow 作為共同記憶體格式的價值:不同資料來源,進入 Query Engine 後可以使用相同的欄式資料結構。
現在可以完整整理成這六步:
1. Data Source 取得或推論 Schema。
2. TableProvider 將 Schema 與資料來源提供給 DataFusion。
3. DataFusion 利用 Schema 驗證 SQL,建立 Logical Plan。
4. Physical Plan 決定如何掃描資料與最佳化讀取。
5. DataSourceExec 在執行時讀取實際資料,
分批產生符合 Schema 的 Arrow RecordBatch。
6. 各個執行節點持續轉換 RecordBatch,
並依運算結果產生新的輸出 Schema。
可以用一句話總結:
Schema 是 DataFusion 理解與規劃資料的依據;
RecordBatch 是資料在 DataFusion 執行流程中實際流動的形式。
今天從資料來源一路追蹤到 Query Engine:
CSV / JSON / Parquet
↓
取得或推論 Schema
↓
TableProvider
↓
Logical Plan → Physical Plan
↓
DataSourceExec 讀取實際資料
↓
Arrow RecordBatch Stream
↓
Execution Operators
↓
Query Result
最重要的觀念是:
Schema 定義資料規格。
TableProvider 提供可查詢的資料來源與 Schema。
DataSourceExec 在執行時讀取資料。
RecordBatch 是資料在 Query Engine 中流動的欄式批次。
明天開始,我們會把視角轉向 SQL 本身:一段 SQL 字串是如何被解析,並逐步變成 DataFusion 的 Logical Plan。