Comece agoraComece grátis

Multiplexando sinks de Parquet

A equipe agora precisa de dois extratos da mesma consulta já limpa: um para empréstimos digitais e outro para físicos. Construa as duas escritas de forma preguiçosa para que o Polars possa planejar a leitura compartilhada uma única vez e, depois, execute tudo em uma única passada.

clean_checkouts já está carregado, assim como DIGITAL_EXPORT_PATH e PHYSICAL_EXPORT_PATH.

Este exercicio faz parte do curso

Dimensionamento e Otimização de Pipelines de Dados com Polars

Ver curso

Instruções do exercicio

  • Construa ambos os sinks de forma preguiçosa para que não executem imediatamente.
  • Execute os dois sinks juntos usando o mecanismo de streaming.

exercicio interativo prático

Tente este exercicio completando este código de exemplo.

# 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 e Executar Código