昨天我們建立了第一個 DataFusion 專案,並使用 SessionContext 執行一段不需要外部資料的 SQL。
今天要開始處理真正的資料檔案。
我們會準備一份簡單的訂單 CSV,將它註冊成 DataFusion 中的 Table,接著查看 Schema,並使用 SELECT、WHERE 和 ORDER BY 查詢資料。
今天的資料流程是:
orders.csv
↓
註冊成 orders Table
↓
DataFusion 建立 Schema
↓
使用 SQL 查詢
↓
輸出結果
延續昨天的 datafusion-day3 專案,在專案根目錄建立 data 資料夾,並在裡面新增 orders.csv:
datafusion-day3
├── Cargo.toml
├── data
│ └── orders.csv
└── src
└── main.rs
orders.csv 內容如下:
order_id,city,amount,is_member
1,Taipei,1200.5,true
2,Taichung,850.0,false
3,Taipei,2300.0,true
4,Kaohsiung,560.5,false
5,Taoyuan,1750.0,true
第一列是欄位名稱,後面每一列代表一筆訂單。
這份 CSV 有四個欄位:
order_id:訂單編號city:城市amount:訂單金額is_member:是否為會員CSV 的全名是 Comma-Separated Values,也就是使用逗號分隔每個欄位值。
它的優點是格式簡單、人類可以直接閱讀,而且幾乎所有資料工具都能處理。
但 CSV 也有一個重要限制:
CSV 檔案本身通常不會完整保存每個欄位的資料型別。
例如,CSV 中的 1200.5 只是一段文字。讀取工具必須判斷它應該是字串、整數還是浮點數。
這就需要用到 Schema。
Schema 用來描述一份資料的結構,包括:
我們預期 orders.csv 的 Schema 大致如下:
order_id → Int64
city → Utf8
amount → Float64
is_member → Boolean
如果沒有正確的 Schema,Query Engine 就無法判斷:
amount >= 1000 是否為合法比較SUM(amount) 是否能執行is_member 是布林值還是字串因此,Schema 是資料進入 Query Engine 時非常重要的一層。
打開 src/main.rs,將程式修改成:
use datafusion::prelude::{CsvReadOptions, SessionContext};
#[tokio::main]
async fn main() -> datafusion::error::Result<()> {
let ctx = SessionContext::new();
let options = CsvReadOptions::new()
.has_header(true);
ctx.register_csv(
"orders",
"data/orders.csv",
options,
)
.await?;
let table = ctx.table("orders").await?;
println!("Inferred schema:\n{:#?}", table.schema());
let dataframe = ctx
.sql(
r#"
SELECT order_id, city, amount, is_member
FROM orders
ORDER BY order_id
"#,
)
.await?;
dataframe.show().await?;
Ok(())
}
執行程式:
cargo run
第一次執行時,終端機會先顯示 DataFusion 推論出的 Schema,接著顯示查詢結果:
+----------+-----------+--------+-----------+
| order_id | city | amount | is_member |
+----------+-----------+--------+-----------+
| 1 | Taipei | 1200.5 | true |
| 2 | Taichung | 850.0 | false |
| 3 | Taipei | 2300.0 | true |
| 4 | Kaohsiung | 560.5 | false |
| 5 | Taoyuan | 1750.0 | true |
+----------+-----------+--------+-----------+
register_csv() 做了什麼?這段程式是今天的核心:
ctx.register_csv(
"orders",
"data/orders.csv",
options,
)
.await?;
三個參數分別代表:
"orders"
└── SQL 中使用的 Table 名稱
"data/orders.csv"
└── CSV 檔案的位置
options
└── 讀取 CSV 時使用的設定
註冊完成後,我們就能在 SQL 中寫:
SELECT *
FROM orders;
這裡要注意,register_csv() 並不是把 CSV 複製進一套新的資料庫。
它比較像是在 SessionContext 中建立一個名稱與資料來源之間的對應:
SQL Table:orders
↓
Data Source:data/orders.csv
當查詢需要讀取 orders 時,DataFusion 才會透過這個對應找到 CSV。
has_header(true)?我們的 CSV 第一列是:
order_id,city,amount,is_member
這一列是欄位名稱,不是訂單資料,因此需要設定:
CsvReadOptions::new().has_header(true)
如果沒有告訴 DataFusion CSV 具有 Header,它可能會把第一列當成一般資料,或使用自動產生的欄位名稱。
CsvReadOptions 還能設定:
例如,有些檔案使用分號分隔:
1;Taipei;1200.5
這時就可以設定:
let options = CsvReadOptions::new()
.has_header(true)
.delimiter(b';');
SELECT 選擇欄位這段 SQL 選擇了四個欄位:
SELECT order_id, city, amount, is_member
FROM orders
ORDER BY order_id;
SELECT 後面列出的欄位,決定查詢結果包含哪些資料。
也可以使用:
SELECT *
FROM orders;
其中 * 代表所有欄位。
雖然 SELECT * 很方便,但在實際專案中,明確列出需要的欄位通常比較容易理解與維護。
例如:
SELECT city, amount
FROM orders;
這個結果只包含城市與訂單金額。
對部分資料格式而言,只選擇需要的欄位也可能減少後續處理量。不過 CSV 是文字格式,讀取時通常仍需要掃描並解析資料列;到了 Parquet,我們會更明顯地看到只讀取需要欄位的優勢。
WHERE 篩選資料接著將 SQL 改成:
let dataframe = ctx
.sql(
r#"
SELECT order_id, city, amount
FROM orders
WHERE amount >= 1000
ORDER BY amount DESC
"#,
)
.await?;
再次執行:
cargo run
結果會只保留金額大於或等於 1000 的訂單:
+----------+---------+--------+
| order_id | city | amount |
+----------+---------+--------+
| 3 | Taipei | 2300.0 |
| 5 | Taoyuan | 1750.0 |
| 1 | Taipei | 1200.5 |
+----------+---------+--------+
這段 SQL 的執行概念是:
TableScan:讀取 orders
↓
Filter:amount >= 1000
↓
Projection:order_id、city、amount
↓
Sort:amount DESC
↓
Result
DESC 代表由大到小排列;如果使用 ASC,則是由小到大排列。
WHERE 也可以使用 AND 或 OR 組合多個條件。
例如,找出金額大於等於 1000,而且是會員的訂單:
SELECT order_id, city, amount
FROM orders
WHERE amount >= 1000
AND is_member = true
ORDER BY amount DESC;
也可以找出來自 Taipei 或 Taoyuan 的訂單:
SELECT order_id, city, amount
FROM orders
WHERE city = 'Taipei'
OR city = 'Taoyuan'
ORDER BY order_id;
要注意:
'Taipei'
1000
true 或 false
這些寫法必須和 Schema 中的資料型別相符。
我們沒有手動指定 Schema,因此 DataFusion 會讀取 CSV 的部分內容並推論型別。
以這份資料來說,它可能推論:
order_id → Int64
city → Utf8
amount → Float64
is_member → Boolean
但 Schema Inference 並不一定永遠正確。
假設某個欄位前面都是數字:
value
100
200
300
N/A
如果推論階段只看到前面的數字,可能將 value 判斷為整數;實際讀到 N/A 時,就可能出現解析錯誤。
常見問題還包括:
在正式的資料 Pipeline 中,如果已經知道資料結構,通常會考慮明確指定 Schema,而不是完全依賴自動推論。
我們也可以使用 Apache Arrow 的 Schema、Field 與 DataType,明確告訴 DataFusion 每個欄位的型別:
use datafusion::arrow::datatypes::{
DataType, Field, Schema,
};
use datafusion::prelude::{
CsvReadOptions, SessionContext,
};
#[tokio::main]
async fn main() -> datafusion::error::Result<()> {
let ctx = SessionContext::new();
let schema = Schema::new(vec![
Field::new("order_id", DataType::Int64, false),
Field::new("city", DataType::Utf8, false),
Field::new("amount", DataType::Float64, false),
Field::new("is_member", DataType::Boolean, false),
]);
let options = CsvReadOptions::new()
.has_header(true)
.schema(&schema);
ctx.register_csv(
"orders",
"data/orders.csv",
options,
)
.await?;
let dataframe = ctx
.sql(
r#"
SELECT order_id, city, amount
FROM orders
WHERE amount >= 1000
ORDER BY amount DESC
"#,
)
.await?;
dataframe.show().await?;
Ok(())
}
Field::new() 的三個參數分別是:
Field::new("amount", DataType::Float64, false)
"amount"
└── 欄位名稱
DataType::Float64
└── 欄位型別
false
└── 不允許 NULL
如果欄位可能為空值,第三個參數可以設為 true。
今天不一定要立刻手動指定 Schema,但要先知道:
CSV 本身缺少完整的型別資訊,因此 Schema 必須由系統推論,或由開發者明確提供。
昨天提到,判斷 SQL 快慢不能只看 SQL 文字,也要知道資料來源如何儲存。
CSV 沒有內建 B/B+ tree 索引,也不像某些欄式格式一樣保存完整的欄位統計資訊。
對一般 CSV 查詢而言,即使 SQL 寫了:
WHERE amount >= 1000
DataFusion 通常仍然需要讀取並解析 CSV 資料,才能判斷哪些資料符合條件。
因此:
WHERE 能減少進入後續 Operator 的資料
但不一定代表:
完全不需要讀取不符合條件的 CSV 內容
等到 Day 6 讀取 Parquet,以及 Day 17~19 討論 Projection Pushdown、Filter Pushdown 與資料裁剪時,就會看到資料格式如何影響 Query Engine 能做的最佳化。
今天完成了第一個真正讀取外部資料的 DataFusion 程式。
我們學到:
register_csv() 可以把 CSV 註冊成 SQL 中可查詢的 Table。has_header(true) 表示 CSV 第一列是欄位名稱。SELECT 決定輸出哪些欄位。WHERE 負責篩選資料。可以先用一句話總結:
註冊 CSV 並不是把檔案搬進資料庫,而是讓 Query Engine 知道一個 Table 名稱對應到哪個 Data Source。
明天將換一種資料格式:
讀取 JSON/NDJSON,觀察半結構化資料如何進入 DataFusion。
SessionContext::register_csv
CsvReadOptions