LoslegenKostenlos starten

Parquet-Sinks multiplexen

Das Team braucht jetzt zwei Exporte aus derselben bereinigten Abfrage: einen für digitale Ausleihen und einen für physische. Erzeuge beide Schreibvorgänge lazy, damit Polars den gemeinsamen Scan einmal planen kann, und führe sie dann in einem Durchlauf zusammen aus.

clean_checkouts ist vorab geladen, ebenso DIGITAL_EXPORT_PATH und PHYSICAL_EXPORT_PATH.

Diese Übung ist Teil des Kurses

<Kurs>Skalieren und Optimieren von Data-Pipelines mit Polars</Kurs>
Kurs ansehen

Übungsanweisungen

  • Erzeuge beide Sinks lazy, damit sie nicht sofort ausgeführt werden.
  • Führe beide Sinks gemeinsam auf der Streaming-Engine aus.

Interaktive praktische Übung

Versuche dich an dieser Übung, indem du diesen Beispielcode vervollständigst.

# 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)
Code bearbeiten und ausführen