CommencerCommencez gratuitement

Multiplexer des sinks Parquet

L'équipe a maintenant besoin de deux extractions à partir de la même requête nettoyée : une pour les emprunts numériques et une pour les emprunts physiques. Construisez les deux écritures en mode paresseux afin que Polars puisse planifier un scan partagé une seule fois, puis exécutez-les ensemble en un seul passage.

clean_checkouts est préchargé, ainsi que DIGITAL_EXPORT_PATH et PHYSICAL_EXPORT_PATH.

Cet exercice fait partie du cours

<cours>Mise à l'échelle et optimisation des pipelines de données avec Polars</cours>
Voir le cours

Instructions de l’exercice

  • Construisez les deux sinks en mode paresseux pour qu'ils ne s'exécutent pas immédiatement.
  • Exécutez les deux sinks ensemble avec le moteur de streaming.

Exercice interactif pratique

Essayez cet exercice en complétant ce code d’exemple.

# 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)
Modifier et exécuter le code