Kom igångKom igång gratis

Multiplexering av Parquet-sinks

Teamet behöver nu två uttag från samma rensade fråga: ett för digitala utlåningar och ett för fysiska. Bygg båda skrivoperationerna lazily så att Polars kan planera den gemensamma skanningen en gång och sedan köra dem tillsammans i ett enda pass.

clean_checkouts är förladdat, tillsammans med DIGITAL_EXPORT_PATH och PHYSICAL_EXPORT_PATH.

Den här övningen är en del av kursen

Skalning och optimering av datapipelines med Polars

Visa kurs

Övningsinstruktioner

  • Bygg båda sinks lazily så att de inte körs direkt.
  • Kör båda sinks tillsammans med streaming-motorn.

Interaktiv övning med praktiskt arbete

Testa den här övningen genom att slutföra den här exempelkoden.

# 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)
Redigera och kör kod