昨天我們將 CSV 註冊成 DataFusion 中的 Table,並使用 SQL 執行欄位選擇與資料篩選。
今天換成另一種常見的資料格式:JSON。
JSON 經常出現在:
與欄位固定、結構平坦的 CSV 相比,JSON 可以包含巢狀物件與陣列,因此常被稱為半結構化資料。
今天的資料流程是:
orders.ndjson
↓
Schema Inference
↓
Arrow Struct / List
↓
註冊成 orders_json Table
↓
使用 SQL 查詢
一般 JSON 可以使用陣列保存多筆資料:
[
{
"order_id": 1,
"city": "Taipei",
"amount": 1200.5
},
{
"order_id": 2,
"city": "Taichung",
"amount": 850.0
}
]
整份檔案是一個 JSON Array,陣列中每個 Object 代表一筆資料。
NDJSON 是 Newline-Delimited JSON 的縮寫。它不使用最外層陣列,而是讓每一行都成為一個完整的 JSON Object:
{"order_id":1,"city":"Taipei","amount":1200.5}
{"order_id":2,"city":"Taichung","amount":850.0}
兩者的差別可以簡化成:
一般 JSON Array
└── 整份檔案是一個 JSON 陣列
NDJSON
└── 每一行都是一筆獨立的 JSON 資料
NDJSON 很適合保存持續產生的事件或 Log,因為系統可以逐行新增與處理資料,不一定要先載入完整的 JSON Array。
要注意的是,NDJSON 的每筆物件必須放在同一行。不能為了排版,把同一個物件拆成多行。
DataFusion 目前可以讀取 NDJSON 與 JSON Array。今天先使用 NDJSON,觀察半結構化資料如何轉換成可以 SQL 查詢的 Table。
延續前幾天的專案,在 data 資料夾新增 orders.ndjson:
datafusion-day3
├── Cargo.toml
├── data
│ ├── orders.csv
│ └── orders.ndjson
└── src
└── main.rs
orders.ndjson 內容如下:
{"order_id":1,"city":"Taipei","amount":1200.5,"customer":{"name":"Amy","level":"gold"},"tags":["mobile","new"]}
{"order_id":2,"city":"Taichung","amount":850.0,"customer":{"name":"Bob","level":"silver"},"tags":["web"]}
{"order_id":3,"city":"Taipei","amount":2300.0,"customer":{"name":"Carol","level":"gold"},"tags":["vip","mobile"]}
{"order_id":4,"city":"Kaohsiung","amount":560.5,"customer":{"name":"David"},"tags":[]}
{"order_id":5,"city":"Taoyuan","amount":1750.0,"customer":{"name":"Eva","level":"silver"},"tags":["web","returning"]}
除了和昨天相同的訂單欄位,這次還加入:
customer:包含顧客名稱與會員等級的巢狀物件tags:保存多個標籤的陣列第四筆資料的 customer 沒有 level,這也是半結構化資料常見的情況:每筆紀錄不一定具有完全相同的欄位。
打開 src/main.rs,改成:
use datafusion::prelude::{
JsonReadOptions, SessionContext,
};
#[tokio::main]
async fn main() -> datafusion::error::Result<()> {
let ctx = SessionContext::new();
let options = JsonReadOptions::default()
.newline_delimited(true)
.file_extension(".ndjson");
ctx.register_json(
"orders_json",
"data/orders.ndjson",
options,
)
.await?;
let table = ctx.table("orders_json").await?;
println!("Inferred schema:\n{:#?}", table.schema());
let dataframe = ctx
.sql(
r#"
SELECT order_id, city, amount
FROM orders_json
ORDER BY order_id
"#,
)
.await?;
dataframe.show().await?;
Ok(())
}
接著執行:
cargo run
應該會看到類似結果:
+----------+-----------+--------+
| order_id | city | amount |
+----------+-----------+--------+
| 1 | Taipei | 1200.5 |
| 2 | Taichung | 850.0 |
| 3 | Taipei | 2300.0 |
| 4 | Kaohsiung | 560.5 |
| 5 | Taoyuan | 1750.0 |
+----------+-----------+--------+
register_json() 做了什麼?今天的核心程式是:
ctx.register_json(
"orders_json",
"data/orders.ndjson",
options,
)
.await?;
它會在 SessionContext 中建立一個名稱與資料來源的對應:
SQL Table:orders_json
↓
Data Source:data/orders.ndjson
註冊完成後,我們就能在 SQL 中使用:
SELECT *
FROM orders_json;
這和昨天的 register_csv() 很相似。它並不是把 JSON 複製進另一套資料庫,而是讓 DataFusion 知道查詢 orders_json 時,應該去哪裡讀取資料。
JsonReadOptions 的設定我們使用了以下設定:
let options = JsonReadOptions::default()
.newline_delimited(true)
.file_extension(".ndjson");
newline_delimited(true) 表示檔案採用 NDJSON 格式,也就是每一行都是獨立的 JSON Object。
.file_extension(".ndjson")
則告訴 DataFusion 這份資料使用 .ndjson 副檔名。
如果讀取的是一般 JSON Array,可以使用:
let options = JsonReadOptions::default()
.newline_delimited(false)
.file_extension(".json");
設定必須和實際檔案格式一致,否則可能出現解析錯誤。
JSON 看起來已經有欄位名稱和值,但 Query Engine 仍然需要 Schema。
DataFusion 需要判斷:
order_id
└── Int64
city
└── Utf8
amount
└── Float64
customer
└── Struct
├── name:Utf8
└── level:Utf8
tags
└── List<Utf8>
其中:
Struct
List
NULL 表示因此,半結構化不代表完全沒有結構。
更精確地說,JSON 的結構可以比 CSV 更有彈性,但當它進入 DataFusion 後,仍然要整理成具有 Schema 的 Arrow 資料,才能交給 Query Engine 執行運算。
我們可以使用欄位存取語法,取得 customer 物件中的資料。
將 SQL 改成:
let dataframe = ctx
.sql(
r#"
SELECT
order_id,
city,
amount,
customer['name'] AS customer_name,
customer['level'] AS customer_level
FROM orders_json
ORDER BY order_id
"#,
)
.await?;
執行後會得到類似結果:
+----------+-----------+--------+---------------+----------------+
| order_id | city | amount | customer_name | customer_level |
+----------+-----------+--------+---------------+----------------+
| 1 | Taipei | 1200.5 | Amy | gold |
| 2 | Taichung | 850.0 | Bob | silver |
| 3 | Taipei | 2300.0 | Carol | gold |
| 4 | Kaohsiung | 560.5 | David | |
| 5 | Taoyuan | 1750.0 | Eva | silver |
+----------+-----------+--------+---------------+----------------+
第四筆資料沒有 customer.level,因此查詢結果會是 NULL。實際顯示方式可能依輸出格式而不同。
這段 SQL:
customer['name']
代表取得 customer Struct 中的 name 欄位。
AS customer_name 則為結果欄位設定較容易閱讀的名稱。
巢狀欄位也可以放進 WHERE 條件。
例如,查詢 Gold 會員的訂單:
SELECT
order_id,
customer['name'] AS customer_name,
city,
amount
FROM orders_json
WHERE customer['level'] = 'gold'
ORDER BY amount DESC;
結果如下:
+----------+---------------+---------+--------+
| order_id | customer_name | city | amount |
+----------+---------------+---------+--------+
| 3 | Carol | Taipei | 2300.0 |
| 1 | Amy | Taipei | 1200.5 |
+----------+---------------+---------+--------+
這表示即使原始資料是具有巢狀結構的 JSON,進入 DataFusion 後,仍然可以使用 SQL 進行欄位選擇、篩選與排序。
JSON 中的:
"tags": ["mobile", "new"]
會被表示成類似 Arrow List<Utf8> 的欄位。
先使用下面的 SQL 查看:
SELECT order_id, tags
FROM orders_json
ORDER BY order_id;
結果中的 tags 仍然是一個 List。
如果希望把陣列中的每個元素展開成獨立資料列,可以使用 UNNEST:
SELECT
order_id,
UNNEST(tags) AS tag
FROM orders_json
ORDER BY order_id;
概念上,原本的資料:
order_id = 1
tags = [mobile, new]
展開後會變成:
order_id | tag
---------+--------
1 | mobile
1 | new
這讓 Query Engine 可以繼續對陣列內容進行篩選、統計與分組。
JSON 允許不同資料列擁有不同欄位,但過度不一致仍然可能造成問題。
例如:
{"amount":1200.5}
{"amount":"unknown"}
第一筆的 amount 是數字,第二筆卻是字串。
DataFusion 在推論 Schema 時,必須決定同一個欄位應該使用哪一種型別。型別不一致可能讓欄位被推論成較寬鬆的型別,也可能在解析或後續運算時發生錯誤。
其他常見問題包括:
所以,處理 JSON 時仍然需要確認資料契約與 Schema,而不是假設任何 JSON 都能直接可靠地查詢。
目前我們已經讀取兩種資料格式:
| 比較項目 | CSV | JSON/NDJSON |
|---|---|---|
| 資料結構 | 通常較平坦 | 可以包含巢狀物件與陣列 |
| 欄位名稱 | 通常來自 Header | 來自 Object Key |
| 型別資訊 | 通常需要推論 | 同樣需要推論 |
| 每筆欄位 | 通常固定 | 可能缺少部分欄位 |
| 人類可讀性 | 高 | 高 |
| 適合情境 | 表格、資料交換 | API、Log、事件與半結構化資料 |
雖然原始格式不同,但進入 DataFusion 後,都需要轉換成 Arrow Schema 與 RecordBatch:
CSV ──────┐
├──→ Arrow RecordBatch → Query Engine
JSON ─────┘
這也是 Query Engine 可以使用相同 SQL Operator 處理不同資料來源的重要原因。
今天使用 DataFusion 讀取 NDJSON,並將半結構化資料註冊成可以 SQL 查詢的 Table。
我們學到:
register_json() 可以將 JSON 資料來源註冊成 Table。NULL 表示。可以先用一句話總結:
JSON 的結構雖然比 CSV 彈性,但進入 Query Engine 後,仍然要轉換成具有明確 Schema 的欄式資料。
明天會介紹第三種資料格式:
Parquet 為什麼特別適合分析型查詢?
JsonReadOptions
SessionContext