多工處理 Parquet sink
團隊現在需要從同一個清理後的查詢產出兩個匯出檔:一個給數位借閱,一個給實體借閱。請把兩個寫出都設為延遲執行,讓 Polars 只規劃並執行一次共享掃描,然後在同一次流程中一併完成。
已預先載入 clean_checkouts,以及 DIGITAL_EXPORT_PATH 和 PHYSICAL_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)