多路复用 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)