Parquet-sinks multiplexen
Het team heeft nu twee exports nodig uit dezelfde opgeschoonde query: één voor digitale uitleningen en één voor fysieke. Bouw beide writes lui zodat Polars de gedeelde scan maar één keer hoeft te plannen, en voer ze dan samen in één keer uit.
clean_checkouts is voorgeladen, net als DIGITAL_EXPORT_PATH en PHYSICAL_EXPORT_PATH.
Deze oefening maakt deel uit van de cursus
Data-pipelines schalen en optimaliseren met Polars
Oefeninstructies
- Bouw beide sinks lui zodat ze niet meteen worden uitgevoerd.
- Voer beide sinks samen uit op de streaming-engine.
Interactieve oefening met praktijkervaring
Probeer deze oefening door deze voorbeeldcode aan te vullen.
# 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)