开始使用免费开始使用

多路复用 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)
编辑并运行代码