將大型查詢直接寫入磁碟

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

Liam Brannigan

Data Scientist & Polars Contributor

巨大輸出問題

示意圖:惰性查詢產生大型、保留列數的抽取結果,應直接寫入磁碟。

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

清理後的請求資料集

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

clean_requests = (
    requests





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

清理後的請求資料集

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

clean_requests = (
    requests
    .with_columns(
        pl.col("CREATED_DATE").dt.date(),
        pl.col("WARD").cast(pl.Int16),
        (pl.col("STATUS") == "Completed").alias("IS_COMPLETED"),
    )
)
使用 Polars 擴充與最佳化資料管線

先收集再寫出

clean_requests.collect(
    engine="streaming"
).write_parquet(
    "311_requests_clean.parquet"
)
使用 Polars 擴充與最佳化資料管線

Sink 到 Parquet

clean_requests.sink_parquet(
    "311_requests_clean.parquet",


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

控制 sink 參數

clean_requests.sink_parquet(
    "311_requests_clean.parquet",
    compression="zstd",
    row_group_size=100_000,
)
使用 Polars 擴充與最佳化資料管線

其他 sink 方法

df.sink_csv("output.csv")
df.sink_ndjson("output.ndjson")
使用 Polars 擴充與最佳化資料管線

分割輸出

clean_requests.sink_parquet(
    pl.PartitionBy(



    ),

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

分割輸出

clean_requests.sink_parquet(
    pl.PartitionBy(
        "311_requests_clean/",
        key="CREATED_DATE",
        max_rows_per_file=1_000_000,
    ),

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

分割輸出

clean_requests.sink_parquet(
    pl.PartitionBy(
        "311_requests_clean/",
        key="CREATED_DATE",
        max_rows_per_file=1_000_000,
    ),
    mkdir=True,
)
使用 Polars 擴充與最佳化資料管線

分割輸出

311_requests_clean/
  CREATED_DATE=2025-12-31/
    00000000.parquet
  CREATED_DATE=2026-01-01/
    00000000.parquet
使用 Polars 擴充與最佳化資料管線

多欄位分割

clean_requests.sink_parquet(
    pl.PartitionBy(
        "311_requests_clean/",
        key=["STATUS", "CREATED_DATE"],
    ),
    mkdir=True,
)
使用 Polars 擴充與最佳化資料管線

建立惰性 sink

internal_sink = clean_requests.sink_parquet(
    "311_requests_internal.parquet",
    lazy=True,
)





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

建立惰性 sink

internal_sink = clean_requests.sink_parquet(
    "311_requests_internal.parquet",
    lazy=True,
)

public_sink = clean_requests.drop("REQUEST_ID","REPORTER").sink_parquet(
    "311_requests_public.parquet",
    lazy=True,
)
使用 Polars 擴充與最佳化資料管線

多工收集 sink

pl.collect_all(
    [internal_sink, public_sink],
)
使用 Polars 擴充與最佳化資料管線

已寫出檔案

311_requests_internal.parquet
311_requests_public.parquet
使用 Polars 擴充與最佳化資料管線

自訂批次 sink

def send_batch(batch: pl.DataFrame) -> None:
    print(batch.height)
    batch_json = batch.write_json()
    # send batch_json to an API
使用 Polars 擴充與最佳化資料管線

自訂批次 sink

clean_requests.sink_batches(
    send_batch,
    chunk_size=50_000,
)
50000
50000
...
29268
使用 Polars 擴充與最佳化資料管線

Sink 工作流程

  • 使用 sink_parquet() 將大型查詢結果直接寫入磁碟
  • 使用 PartitionBy() 將輸出切成分割資料集
  • 用惰性 sink 搭配 collect_all() 規劃並一次執行多個輸出
  • 當每個批次要傳給自訂函式時,使用 sink_batches()

資料分流系統的等角視圖,將資料分配到有組織的分割抽屜

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

一起來練習吧!

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

Preparing Video For Download...