CommencezCommencez gratuitement

Multiplexage de « sinks » Parquet

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

clean_checkouts est préchargé, tout comme DIGITAL_EXPORT_PATH et PHYSICAL_EXPORT_PATH.

Cette activité fait partie du cours

Mise à l'échelle et optimisation des pipelines de données avec Polars

Voir le cours

Instructions de l’exercice

  • Créez les deux « sinks » en mode paresseux afin qu'ils ne s'exécutent pas immédiatement.
  • Exécutez les deux « sinks » ensemble avec le moteur de diffusion en continu (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