Začněte nyníZačněte zdarma

Multiplexování Parquet sinků

Tým teď potřebuje dva výpisy ze stejného vyčištěného dotazu: jeden pro digitální výpůjčky a druhý pro fyzické. Sestav oba zápisy líně, aby Polars naplánoval sdílený scan dat jen jednou a pak je spustil v jediném průchodu.

clean_checkouts je přednahrané, stejně jako DIGITAL_EXPORT_PATH a PHYSICAL_EXPORT_PATH.

Toto cvičení je součástí kurzu

Scaling and Optimizing Data Pipelines with Polars

Zobrazit kurz

Pokyny k cvičení

  • Sestav oba sinky líně, aby se nespustily okamžitě.
  • Spusť oba sinky zároveň pomocí streaming enginu.

Interaktivní cvičení na vyzkoušení si v praxi

Vyzkoušejte si toto cvičení dokončením tohoto ukázkového kódu.

# 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)
Upravit a spustit kód