EmpezarEmpieza gratis

Multiplexar sinks Parquet

El equipo ahora necesita dos extractos de la misma consulta ya depurada: uno para préstamos digitales y otro para físicos. Construye ambas escrituras de forma perezosa para que Polars pueda planificar el escaneo compartido una sola vez, y luego ejecútalas juntas en una pasada.

clean_checkouts está precargado, junto con DIGITAL_EXPORT_PATH y PHYSICAL_EXPORT_PATH.

Este ejercicio forma parte del curso

Escala y optimiza canalizaciones de datos con Polars

Ver curso

Instrucciones del ejercicio

  • Construye ambos sinks de forma perezosa para que no se ejecuten inmediatamente.
  • Ejecuta ambos sinks juntos con el motor de streaming.

ejercicio interactivo práctico

Prueba este ejercicio completando este código de ejemplo.

# 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)
Editar y ejecutar código