Zacznij terazZacznij za darmo

Multipleksowanie ujść Parquet

Zespół potrzebuje teraz dwóch wyciągów z tego samego oczyszczonego zapytania: jednego dla wypożyczeń cyfrowych i jednego dla fizycznych. Zbuduj oba zapisy leniwie, żeby Polars mógł zaplanować wspólny skan raz i wykonać go w jednym przebiegu.

clean_checkouts jest wstępnie załadowane, podobnie jak DIGITAL_EXPORT_PATH i PHYSICAL_EXPORT_PATH.

To ćwiczenie jest częścią kursu

Skalowanie i optymalizacja potoków danych w Polars

Zobacz kurs

Instrukcje do ćwiczenia

  • Zbuduj oba ujścia leniwie, tak żeby nie wykonywały się od razu.
  • Uruchom oba ujścia jednocześnie na silniku strumieniowym.

Interaktywne ćwiczenie praktyczne

Spróbuj tego ćwiczenia, uzupełniając ten przykładowy kod.

# 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)
Edytuj i uruchom kod