iT邦幫忙

2026 iThome 鐵人賽

DAY 10
0

昨天有讀者問:

DataFusion 是如何利用 Schema 產生 RecordBatch?DataFusion 和 Schema 的關係到底是什麼?

簡短的答案是:

Schema 不會憑空產生 RecordBatch。

Schema 定義資料規格;
Data Source 讀取實際資料;
DataFusion 依照 Schema 規劃與驗證查詢,
並在執行時讓資料來源分批產生符合規格的 RecordBatch。

也就是說,Schema 像資料的設計圖,RecordBatch 則是依照設計圖裝好的實際資料。

今天就從 CSV、JSON、Parquet 等資料來源出發,追蹤它們如何進入 DataFusion,最後成為 Query Engine 能處理的 Arrow RecordBatch。


回顧:DataFusion 實際處理的是 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。


從資料來源到 RecordBatch 的流程

DataFusion 從資料來源讀取資料、產生 Arrow RecordBatch,並執行查詢的簡化流程

圖 1:DataFusion 從資料來源讀取資料、產生 Arrow RecordBatch,並執行查詢的簡化流程。圖由本文整理。

這張圖刻意分成兩個階段:

規劃階段
→ 先確認資料結構,決定 SQL 要如何執行。

執行階段
→ 實際讀取資料,分批產生 RecordBatch 並完成運算。

這個區分很重要。DataFusion 不會一收到 SQL 就直接把所有檔案讀進記憶體。


規劃階段:DataFusion 先理解資料結構

假設我們執行:

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 的方式不太一樣。

資料格式 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 是什麼?

可以把 TableProvider 想成 DataFusion 與資料來源之間的橋樑。

它負責告訴 DataFusion:

這張表的 Schema 是什麼?
資料在哪裡?
能否裁剪欄位?
能否下推篩選條件?
實際掃描資料時,應該建立什麼執行計畫?

因此,SQL 裡的:

FROM orders

不是單純指向一個 orders.parquet 檔案,而是指向一個可查詢的資料表來源。

DataFusion 內建的 CSV、JSON、Parquet 讀取支援,會將檔案來源包裝成可查詢的 TableProvider。


Physical Plan:決定真正怎麼讀資料

Logical Plan 完成後,DataFusion 會建立 Physical Plan。

它決定更接近實作層面的事情,例如:

要讀取哪些檔案?
要讀取哪些欄位?
能否跳過部分資料?
資料要如何分區與平行處理?

以上面的 SQL 為例,order_id 完全沒有被使用。

因此,Physical Plan 可以只讀取:

city
amount
is_member

這就是 Day 7 提過的 Column Pruning。

若資料來源是 Parquet,DataFusion 還可能依條件與檔案統計資訊,跳過不可能符合條件的資料區塊。

這個階段仍是在規劃「如何讀取」;真正的資料 I/O 會在執行階段發生。


執行階段:DataSourceExec 讀取實際資料

開始執行查詢後,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 如何一路流過 Query Engine?

接著,這批資料會流過不同的執行節點。

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 最後為什麼能共用同一套 Query Engine?

三種格式的原始資料差異很大:

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 後可以使用相同的欄式資料結構。


回答讀者問題:DataFusion 與 Schema 的關係

現在可以完整整理成這六步:

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。


延伸閱讀


上一篇
Day 9|Schema、Array 與 RecordBatch
下一篇
Day 11|SQL 不只是字串:Query Engine 如何理解 SQL?
系列文
深入 SQL 查詢引擎:30 天用 Rust 與 Apache DataFusion 解構資料處理流程13
圖片
  熱門推薦
圖片
{{ item.channelVendor }} | {{ item.webinarstarted }} |
{{ formatDate(item.duration) }}
直播中

尚未有邦友留言

立即登入留言