使用串流與 GPU 執行引擎

使用 Polars 擴充與最佳化資料管線

Liam Brannigan

Data Scientist & Polars Contributor

查詢執行引擎

流程圖:比較查詢從 API 經查詢最佳化器到記憶體引擎。

使用 Polars 擴充與最佳化資料管線

查詢執行引擎

流程圖:比較查詢從 API 經查詢最佳化器到記憶體引擎。

使用 Polars 擴充與最佳化資料管線

查詢執行引擎

流程圖:比較查詢從 API 經查詢最佳化器到記憶體與串流引擎。

使用 Polars 擴充與最佳化資料管線

查詢執行引擎

流程圖:比較查詢從 API 經查詢最佳化器到記憶體、串流與 GPU 引擎。

使用 Polars 擴充與最佳化資料管線

惰性請求摘要

requests = pl.scan_parquet("311_Service_Requests.parquet")

query = (
    requests
    .filter(pl.col("STATUS") == "Completed")
    .group_by("DEPARTMENT")
    .len()
    .sort("len", descending=True)
)
使用 Polars 擴充與最佳化資料管線

使用預設引擎

query.collect()
shape: (5, 2)
| DEPARTMENT        | len     |
| ---               | ---     |
| str               | u32     |
|-------------------|---------|
| 311 City Services | 4859161 |
| Sanitation        | 3406631 |
| Aviation          | 2337842 |
| Transportation    | 1468240 |
| Water             | 416535  |
使用 Polars 擴充與最佳化資料管線

指定串流引擎

query.collect(
    engine="streaming"
)
shape: (5, 2)
| DEPARTMENT        | len     |
| ---               | ---     |
| str               | u32     |
|-------------------|---------|
| 311 City Services | 4859161 |
| Sanitation        | 3406631 |
| Aviation          | 2337842 |
| Transportation    | 1468240 |
| Water             | 416535  |
使用 Polars 擴充與最佳化資料管線

指定串流引擎

大型資料表以串流處理的圖片。

  • 不支援的串流步驟會自動回退到預設引擎
使用 Polars 擴充與最佳化資料管線

指定串流引擎

import polars as pl
pl.Config.set_engine_affinity("streaming")
使用 Polars 擴充與最佳化資料管線

指定 GPU 引擎

query.collect(
    engine="gpu"
)

$$

  • 需要 NVIDIA GPU 硬體
使用 Polars 擴充與最佳化資料管線

選擇引擎

$$

  • 先用預設引擎;僅在記憶體受限時再切換

$$

  • 串流適用於超大、超出記憶體的查詢

$$

  • GPU 追求速度,但僅在有 NVIDIA 硬體時使用

選擇引擎

使用 Polars 擴充與最佳化資料管線

分批處理

completed_requests = (
    requests
    .filter(pl.col("STATUS") == "Completed")
    .select("SR_NUMBER", "TYPE", "DEPARTMENT", "CREATED_DATE")
)
使用 Polars 擴充與最佳化資料管線

處理每個批次

for batch in completed_requests.collect_batches(

):


使用 Polars 擴充與最佳化資料管線

處理每個批次

for batch in completed_requests.collect_batches(
    chunk_size=50_000,
):


使用 Polars 擴充與最佳化資料管線

處理每個批次

for batch in completed_requests.collect_batches(
    chunk_size=50_000,
):
    print(batch.shape)
    batch_json = batch.write_json()
(50000, 4)
(50000, 4)
...
(50000, 4)
(12409, 4)
使用 Polars 擴充與最佳化資料管線

一起來練習吧!

使用 Polars 擴充與最佳化資料管線

Preparing Video For Download...