iT邦幫忙

2026 iThome 鐵人賽

DAY 4
1

昨天我們建立了第一個 DataFusion 專案,並使用 SessionContext 執行一段不需要外部資料的 SQL。

今天要開始處理真正的資料檔案。

我們會準備一份簡單的訂單 CSV,將它註冊成 DataFusion 中的 Table,接著查看 Schema,並使用 SELECTWHEREORDER BY 查詢資料。

今天的資料流程是:

orders.csv
    ↓
註冊成 orders Table
    ↓
DataFusion 建立 Schema
    ↓
使用 SQL 查詢
    ↓
輸出結果

準備 CSV 資料

延續昨天的 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?

Schema 用來描述一份資料的結構,包括:

  • 有哪些欄位
  • 欄位名稱是什麼
  • 每個欄位的資料型別
  • 欄位是否允許空值

我們預期 orders.csv 的 Schema 大致如下:

order_id   → Int64
city       → Utf8
amount     → Float64
is_member  → Boolean

如果沒有正確的 Schema,Query Engine 就無法判斷:

  • amount >= 1000 是否為合法比較
  • SUM(amount) 是否能執行
  • is_member 是布林值還是字串
  • 不同 Table 的欄位能不能 JOIN

因此,Schema 是資料進入 Query Engine 時非常重要的一層。


將 CSV 註冊成 Table

打開 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 還能設定:

  • 分隔符號
  • 引號字元
  • Escape 字元
  • Comment 字元
  • Schema 推論筆數
  • 壓縮格式
  • 空值規則

例如,有些檔案使用分號分隔:

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 也可以使用 ANDOR 組合多個條件。

例如,找出金額大於等於 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
  • 布林值可以使用 truefalse

這些寫法必須和 Schema 中的資料型別相符。


CSV 的 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,而不是完全依賴自動推論。


手動指定 Schema

我們也可以使用 Apache Arrow 的 SchemaFieldDataType,明確告訴 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 必須由系統推論,或由開發者明確提供。


CSV 與昨天的效能問題

昨天提到,判斷 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 程式。

我們學到:

  1. CSV 使用逗號分隔欄位,但本身缺少完整的型別資訊。
  2. Schema 描述欄位名稱、資料型別與是否允許空值。
  3. register_csv() 可以把 CSV 註冊成 SQL 中可查詢的 Table。
  4. has_header(true) 表示 CSV 第一列是欄位名稱。
  5. SELECT 決定輸出哪些欄位。
  6. WHERE 負責篩選資料。
  7. DataFusion 可以推論 CSV Schema,也可以由開發者手動指定。
  8. CSV 的儲存特性會限制 Query Engine 可以進行的資料裁剪方式。

可以先用一句話總結:

註冊 CSV 並不是把檔案搬進資料庫,而是讓 Query Engine 知道一個 Table 名稱對應到哪個 Data Source。

明天將換一種資料格式:

讀取 JSON/NDJSON,觀察半結構化資料如何進入 DataFusion。


延伸閱讀


上一篇
Day 3|建立第一個 DataFusion 專案:相同結果的 SQL,真的一樣快嗎?
下一篇
Day 5|JSON 資料也能直接 SQL 查詢
系列文
深入 SQL 查詢引擎:30 天用 Rust 與 Apache DataFusion 解構資料處理流程5
圖片
  熱門推薦
圖片
{{ item.channelVendor }} | {{ item.webinarstarted }} |
{{ formatDate(item.duration) }}
直播中

尚未有邦友留言

立即登入留言