ÎncepețiÎncepe gratuit

Multiplexarea sink-urilor Parquet

Echipa are nevoie acum de două extrase din același query curățat: unul pentru împrumuturile digitale și altul pentru cele fizice. Construiește ambele scrieri în mod lazy, astfel încât Polars să poată planifica o singură dată scanarea comună, apoi să le execute împreună într-o singură trecere.

clean_checkouts este preîncărcat, împreună cu DIGITAL_EXPORT_PATH și PHYSICAL_EXPORT_PATH.

Acest exercițiu face parte din cursul

Scalarea și optimizarea pipeline-urilor de date cu Polars

Vezi cursul

Instrucțiuni pentru exercițiu

  • Construiește ambele sink-uri în mod lazy, astfel încât să nu se execute imediat.
  • Rulează ambele sink-uri împreună pe motorul de streaming.

Exercițiu interactiv practic

Încearcă acest exercițiu completând acest cod de exemplu.

# 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)
Editează și rulează codul