Skriva stora frågor direkt till disk

Skalning och optimering av datapipelines med Polars

Liam Brannigan

Data Scientist & Polars Contributor

Problemet med stora utdata

Diagram som visar en lat fråga som producerar ett stort radbevarande extrakt som ska skrivas direkt till disk.

Skalning och optimering av datapipelines med Polars

En rensad begärandedatamängd

requests = pl.scan_parquet("311_Service_Requests.parquet")

clean_requests = (
    requests





)
Skalning och optimering av datapipelines med Polars

En rensad begärandedatamängd

requests = pl.scan_parquet("311_Service_Requests.parquet")

clean_requests = (
    requests
    .with_columns(
        pl.col("CREATED_DATE").dt.date(),
        pl.col("WARD").cast(pl.Int16),
        (pl.col("STATUS") == "Completed").alias("IS_COMPLETED"),
    )
)
Skalning och optimering av datapipelines med Polars

Samla in och skriv

clean_requests.collect(
    engine="streaming"
).write_parquet(
    "311_requests_clean.parquet"
)
Skalning och optimering av datapipelines med Polars

Skriva till Parquet med sink

clean_requests.sink_parquet(
    "311_requests_clean.parquet",


)
Skalning och optimering av datapipelines med Polars

Styra sink-inställningarna

clean_requests.sink_parquet(
    "311_requests_clean.parquet",
    compression="zstd",
    row_group_size=100_000,
)
Skalning och optimering av datapipelines med Polars

Andra sink-metoder

df.sink_csv("output.csv")
df.sink_ndjson("output.ndjson")
Skalning och optimering av datapipelines med Polars

Partitionerade utdata

clean_requests.sink_parquet(
    pl.PartitionBy(



    ),

)
Skalning och optimering av datapipelines med Polars

Partitionerade utdata

clean_requests.sink_parquet(
    pl.PartitionBy(
        "311_requests_clean/",
        key="CREATED_DATE",
        max_rows_per_file=1_000_000,
    ),

)
Skalning och optimering av datapipelines med Polars

Partitionerade utdata

clean_requests.sink_parquet(
    pl.PartitionBy(
        "311_requests_clean/",
        key="CREATED_DATE",
        max_rows_per_file=1_000_000,
    ),
    mkdir=True,
)
Skalning och optimering av datapipelines med Polars

Partitionerade utdata

311_requests_clean/
  CREATED_DATE=2025-12-31/
    00000000.parquet
  CREATED_DATE=2026-01-01/
    00000000.parquet
Skalning och optimering av datapipelines med Polars

Partitioner med flera kolumner

clean_requests.sink_parquet(
    pl.PartitionBy(
        "311_requests_clean/",
        key=["STATUS", "CREATED_DATE"],
    ),
    mkdir=True,
)
Skalning och optimering av datapipelines med Polars

Skapa lata sinks

internal_sink = clean_requests.sink_parquet(
    "311_requests_internal.parquet",
    lazy=True,
)





Skalning och optimering av datapipelines med Polars

Skapa lata sinks

internal_sink = clean_requests.sink_parquet(
    "311_requests_internal.parquet",
    lazy=True,
)

public_sink = clean_requests.drop("REQUEST_ID","REPORTER").sink_parquet(
    "311_requests_public.parquet",
    lazy=True,
)
Skalning och optimering av datapipelines med Polars

Multiplexering av sinks

pl.collect_all(
    [internal_sink, public_sink],
)
Skalning och optimering av datapipelines med Polars

Skrivna utdatafiler

311_requests_internal.parquet
311_requests_public.parquet
Skalning och optimering av datapipelines med Polars

Anpassade batch-sinks

def send_batch(batch: pl.DataFrame) -> None:
    print(batch.height)
    batch_json = batch.write_json()
    # send batch_json to an API
Skalning och optimering av datapipelines med Polars

Anpassade batch-sinks

clean_requests.sink_batches(
    send_batch,
    chunk_size=50_000,
)
50000
50000
...
29268
Skalning och optimering av datapipelines med Polars

Sink-arbetsflöde

  • Använd sink_parquet() för att skriva stora frågeresultat direkt till disk
  • Använd PartitionBy() för att dela upp utdata i en partitionerad datamängd
  • Använd lata sinks med collect_all() för att planera och köra flera utdatafiler tillsammans
  • Använd sink_batches() när varje batch ska skickas till en anpassad funktion

Isometrisk vy av ett datasorteringssystem som fördelar data i organiserade partitionerade lådor

Skalning och optimering av datapipelines med Polars

Nu kör vi en övning!

Skalning och optimering av datapipelines med Polars

Preparing Video For Download...