前兩天,我們先從概念上認識了 Query Engine,也知道 DataFusion 是一套使用 Rust 與 Apache Arrow 建立的查詢引擎。
今天終於要開始寫程式了!
我們會完成三件事:
SessionContext 執行第一段 SQL最後再做一個小實驗:
兩段結果相同的 SQL,執行速度也會一樣嗎?
開始前,可以先確認電腦已經安裝 Rust:
rustc --version
cargo --version
如果兩個指令都有顯示版本,就可以建立專案:
cargo new datafusion-day3
cd datafusion-day3
專案結構大致如下:
datafusion-day3
├── Cargo.toml
└── src
└── main.rs
其中:
Cargo.toml:管理專案資訊與相依套件src/main.rs:程式進入點截至本文撰寫時,使用的 DataFusion 版本為 54.1.0。
在專案目錄執行:
cargo add datafusion@54.1.0
cargo add tokio --features macros,rt-multi-thread
為什麼除了 DataFusion,還要加入 Tokio?
DataFusion 的查詢、檔案讀取與執行 API 大量使用 Rust 的 async/await。Tokio 提供執行非同步程式需要的 Runtime。
完成後,Cargo.toml 會出現類似設定:
[dependencies]
datafusion = "54.1.0"
tokio = { version = "1.53.1", features = ["macros", "rt-multi-thread"] }
第一次編譯 DataFusion 需要下載及編譯不少相依套件,因此可能會等待一段時間,之後再次編譯通常會快很多。
打開 src/main.rs,將內容改成:
use datafusion::prelude::SessionContext;
#[tokio::main]
async fn main() -> datafusion::error::Result<()> {
let ctx = SessionContext::new();
let sql = r#"
SELECT
n,
n * n AS square
FROM (
VALUES (1), (2), (3), (4), (5)
) AS numbers(n)
WHERE n > 2
ORDER BY n
"#;
let dataframe = ctx.sql(sql).await?;
dataframe.show().await?;
Ok(())
}
接著執行:
cargo run
應該會看到類似結果:
+---+--------+
| n | square |
+---+--------+
| 3 | 9 |
| 4 | 16 |
| 5 | 25 |
+---+--------+
我們沒有讀取外部檔案,而是先使用 SQL 的 VALUES 建立一組簡單資料。
這段 SQL 做了幾件事:
VALUES 建立 1 到 5 的資料WHERE n > 2 保留 3、4、5n * n 計算平方ORDER BY n 排序結果SessionContext 是什麼?程式中最重要的角色是:
let ctx = SessionContext::new();
SessionContext 是使用 DataFusion 執行查詢時的主要入口。
它負責管理一個查詢 Session 所需要的狀態,例如:
之後讀取 CSV、JSON、Parquet,以及建立 UDF 時,也會繼續使用它。
接著:
let dataframe = ctx.sql(sql).await?;
這行會接收 SQL,並建立代表查詢的 DataFrame。此時 DataFusion 會進行 SQL 解析與查詢規劃,但不一定已經完成所有資料運算。
真正要求 DataFusion 執行並顯示結果的是:
dataframe.show().await?;
可以先簡化理解成:
ctx.sql()
↓
解析 SQL 並建立查詢計畫
↓
DataFrame
show()
↓
建立並執行 Physical Plan
↓
顯示結果
DataFusion 採用 Lazy Evaluation。也就是先描述要執行的查詢,等到呼叫 show()、collect() 等方法時,才真正觸發執行。
有人建議我可以比較功能相同的 SQL,並分析不同寫法的執行速度。
這是一個很有趣的問題,但在比較以前,要先區分三個角色:
Rust Compiler
└── 將 Rust 程式編譯成可執行檔
DataFusion Query Optimizer
└── 最佳化 SQL 產生的 Query Plan
Data Source
└── 決定資料如何儲存及讀取
即使我們假設 Rust Compiler 沒有額外最佳化,DataFusion 的 Query Optimizer 仍然可能改寫 Query Plan。
因此,不能只看 SQL 字串長短,就直接判斷哪一段比較快。
AND TRUE 會比較慢嗎?先看兩段 SQL。
SELECT n, n * n AS square
FROM (
VALUES (1), (2), (3), (4), (5)
) AS numbers(n)
WHERE n > 2;
SELECT n, n * n AS square
FROM (
VALUES (1), (2), (3), (4), (5)
) AS numbers(n)
WHERE n > 2 AND TRUE;
兩段 SQL 的結果完全相同。
如果完全按照文字逐步運算,SQL B 好像多判斷了一次 TRUE,理論上需要多做一點工作。
但是 DataFusion 的 Query Optimizer 可以進行 Expression Simplification,把:
n > 2 AND TRUE
簡化成:
n > 2
經過最佳化後,兩段 SQL 可能產生相同或非常接近的 Physical Plan,因此實際效能不一定有差別。
這告訴我們:
SQL 看起來比較長,不代表執行時一定比較慢。
雖然 Query Plan 會在 Day 13 進一步介紹,今天可以先使用 EXPLAIN 偷看一下。
把查詢改成:
let sql = r#"
EXPLAIN VERBOSE FORMAT INDENT
SELECT n, n * n AS square
FROM (
VALUES (1), (2), (3), (4), (5)
) AS numbers(n)
WHERE n > 2 AND TRUE
"#;
再執行:
cargo run
輸出會包含不同階段的 Logical Plan 與 Physical Plan。
內容可能很多,今天不需要全部看懂。可以先尋找 simplify_expressions,觀察 AND TRUE 是否在最佳化過程中被移除。
UNION 和 UNION ALL再看一組在特定條件下結果相同的 SQL。
假設資料中的數字不重複,而且兩個查詢條件完全不重疊:
UNIONSELECT n
FROM (
VALUES (1), (2), (3), (4), (5)
) AS numbers(n)
WHERE n <= 3
UNION
SELECT n
FROM (
VALUES (1), (2), (3), (4), (5)
) AS numbers(n)
WHERE n > 3;
UNION ALLSELECT n
FROM (
VALUES (1), (2), (3), (4), (5)
) AS numbers(n)
WHERE n <= 3
UNION ALL
SELECT n
FROM (
VALUES (1), (2), (3), (4), (5)
) AS numbers(n)
WHERE n > 3;
在這份資料與條件下,兩者會得到相同資料。
但一般來說:
UNION:需要移除重複資料UNION ALL:直接合併兩邊的結果如果我們已經確定兩邊不可能產生重複資料,UNION ALL 通常能避免去重所需的額外工作。
不過,這個結論有一個重要前提:
只有在可以保留重複資料,或能確定資料不會重複時,才能用
UNION ALL取代UNION。
如果資料可能重複,兩段 SQL 的結果就不相同,不能只為了效能直接替換。
在傳統關聯式資料庫中,B-tree 或 B+ tree 經常被用來建立索引。
例如:
SELECT *
FROM orders
WHERE order_id = 100;
如果 order_id 有合適的索引,資料庫可能不需要掃描整張表,而是透過索引更快找到資料。
但 DataFusion 的定位是可嵌入的 Query Engine,不是內建固定儲存引擎的完整資料庫。它不會假設所有資料都儲存在 B/B+ tree 中。
DataFusion 可能處理:
因此,在 DataFusion 中討論查詢效能時,需要先知道資料來源。
例如:
Arrow RecordBatch
└── 主要在記憶體中以欄式 Array 處理
CSV
└── 通常需要解析文字內容
Parquet
└── 可以利用欄位資訊、Metadata 與 Statistics 減少讀取
自訂 Data Source
└── 可能擁有自己的索引或查詢能力
所以「哪段 SQL 比較快」不能只根據 SQL 文字判斷,還要考慮:
CSV、Parquet、Filter Pushdown 與資料裁剪,會在後續文章中實際比較。
今天的 VALUES 只有五筆資料。如果直接測量兩段 SQL 的執行時間,結果很容易受到以下因素影響:
因此,今天的重點不是宣稱其中一段快了幾毫秒,而是先學會觀察:
SQL
↓
Logical Plan
↓
Optimizer
↓
Physical Plan
↓
Execution
如果未來真的要做效能比較,應該使用足夠大的資料、固定環境、重複執行,並使用 Benchmark 工具,而不是只執行一次 cargo run 就下結論。
今天完成了第一個 DataFusion 專案,並使用 SessionContext 執行第一段 SQL。
我們也發現:
SessionContext 是使用 DataFusion 的主要入口。ctx.sql() 會建立查詢,show() 會觸發執行。可以先用一句話總結:
我們寫的是 SQL,但 Query Engine 真正執行的是最佳化後的 Physical Plan。
明天會開始載入真正的資料:
從 CSV 開始,將一份檔案註冊成可以使用 SQL 查詢的 Table。