開始使用免費開始

多工處理 Parquet sink

團隊現在需要從同一個清理後的查詢產出兩個匯出檔:一個給數位借閱,一個給實體借閱。請把兩個寫出都設為延遲執行,讓 Polars 只規劃並執行一次共享掃描,然後在同一次流程中一併完成。

已預先載入 clean_checkouts,以及 DIGITAL_EXPORT_PATHPHYSICAL_EXPORT_PATH

本練習屬於課程

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

檢視課程

練習說明

  • 以延遲方式建立兩個 sink,避免馬上執行。
  • 在 streaming 引擎上一併執行這兩個 sink。

動手互動練習

試著完成這個範例程式碼,體驗一下這個練習。

# Lazy sink for digital
digital_sink = (
    clean_checkouts
    .filter(pl.col("use") == "Digital")
    .sink_parquet(DIGITAL_EXPORT_PATH, lazy=____)
)

# Lazy sink for physical
physical_sink = (
    clean_checkouts
    .filter(pl.col("use") == "Physical")
    .sink_parquet(PHYSICAL_EXPORT_PATH, lazy=____)
)

# Run both sinks together
pl.____([digital_sink, physical_sink], engine="streaming")

# Check the row counts of each extract
result = pl.DataFrame(
    {
        "extract": ["digital", "physical"],
        "rows": [
            pl.scan_parquet(DIGITAL_EXPORT_PATH).select(pl.len()).collect().item(),
            pl.scan_parquet(PHYSICAL_EXPORT_PATH).select(pl.len()).collect().item(),
        ],
    }
)
print(result)
編輯並執行程式碼