Ukládání velkých dotazů na disk

Scaling and Optimizing Data Pipelines with Polars

Liam Brannigan

Data Scientist & Polars Contributor

Problém s velkým výstupem

Diagram znázorňující líný dotaz produkující velký výstup, který má být zapsán přímo na disk.

Scaling and Optimizing Data Pipelines with Polars

Vyčištěná datová sada požadavků

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

clean_requests = (
    requests





)
Scaling and Optimizing Data Pipelines with Polars

Vyčištěná datová sada požadavků

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"),
    )
)
Scaling and Optimizing Data Pipelines with Polars

Shromáždění a zápis

clean_requests.collect(
    engine="streaming"
).write_parquet(
    "311_requests_clean.parquet"
)
Scaling and Optimizing Data Pipelines with Polars

Zápis do Parquet pomocí sinku

clean_requests.sink_parquet(
    "311_requests_clean.parquet",


)
Scaling and Optimizing Data Pipelines with Polars

Nastavení sinku

clean_requests.sink_parquet(
    "311_requests_clean.parquet",
    compression="zstd",
    row_group_size=100_000,
)
Scaling and Optimizing Data Pipelines with Polars

Další metody sinku

df.sink_csv("output.csv")
df.sink_ndjson("output.ndjson")
Scaling and Optimizing Data Pipelines with Polars

Rozdělený výstup

clean_requests.sink_parquet(
    pl.PartitionBy(



    ),

)
Scaling and Optimizing Data Pipelines with Polars

Rozdělený výstup

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

)
Scaling and Optimizing Data Pipelines with Polars

Rozdělený výstup

clean_requests.sink_parquet(
    pl.PartitionBy(
        "311_requests_clean/",
        key="CREATED_DATE",
        max_rows_per_file=1_000_000,
    ),
    mkdir=True,
)
Scaling and Optimizing Data Pipelines with Polars

Rozdělený výstup

311_requests_clean/
  CREATED_DATE=2025-12-31/
    00000000.parquet
  CREATED_DATE=2026-01-01/
    00000000.parquet
Scaling and Optimizing Data Pipelines with Polars

Rozdělení podle více sloupců

clean_requests.sink_parquet(
    pl.PartitionBy(
        "311_requests_clean/",
        key=["STATUS", "CREATED_DATE"],
    ),
    mkdir=True,
)
Scaling and Optimizing Data Pipelines with Polars

Vytváření líných sinků

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





Scaling and Optimizing Data Pipelines with Polars

Vytváření líných sinků

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,
)
Scaling and Optimizing Data Pipelines with Polars

Multiplexování sinků

pl.collect_all(
    [internal_sink, public_sink],
)
Scaling and Optimizing Data Pipelines with Polars

Zapsané výstupy

311_requests_internal.parquet
311_requests_public.parquet
Scaling and Optimizing Data Pipelines with Polars

Vlastní dávkové sinky

def send_batch(batch: pl.DataFrame) -> None:
    print(batch.height)
    batch_json = batch.write_json()
    # send batch_json to an API
Scaling and Optimizing Data Pipelines with Polars

Vlastní dávkové sinky

clean_requests.sink_batches(
    send_batch,
    chunk_size=50_000,
)
50000
50000
...
29268
Scaling and Optimizing Data Pipelines with Polars

Pracovní postup sinku

  • Použijte sink_parquet() k zápisu velkých výsledků přímo na disk
  • Použijte PartitionBy() k rozdělení výstupu do oddílů
  • Použijte líné sinky s collect_all() pro plánování a spuštění více výstupů najednou
  • Použijte sink_batches(), pokud má každá dávka jít do vlastní funkce

Izometrický pohled na systém třídění dat rozdělující data do organizovaných oddílů

Scaling and Optimizing Data Pipelines with Polars

Lass uns üben!

Scaling and Optimizing Data Pipelines with Polars

Preparing Video For Download...