Multiplexing dei sink Parquet
Il team ora ha bisogno di due estrazioni dalla stessa query pulita: una per i prestiti digitali e una per quelli fisici. Crea entrambe le scritture in modo lazy così che Polars possa pianificare una sola scansione condivisa, poi eseguirle insieme in un unico passaggio.
clean_checkouts è già caricato, insieme a DIGITAL_EXPORT_PATH e PHYSICAL_EXPORT_PATH.
Questo esercizio fa parte del corso
Scalare e ottimizzare le pipeline di dati con Polars
Istruzioni dell'esercizio
- Crea entrambi i sink in modo lazy così da non eseguirli immediatamente.
- Esegui entrambi i sink insieme con il motore di streaming.
esercizio interattivo pratico
Prova questo esercizio completando questo codice di esempio.
# 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)