Menulis kueri besar ke disk

Menskalakan dan Mengoptimalkan Pipeline Data dengan Polars

Liam Brannigan

Data Scientist & Polars Contributor

Masalah output besar

Diagram yang menunjukkan kueri lazy menghasilkan ekstrak besar yang mempertahankan baris dan sebaiknya ditulis langsung ke disk.

Menskalakan dan Mengoptimalkan Pipeline Data dengan Polars

Dataset permintaan yang dibersihkan

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

clean_requests = (
    requests





)
Menskalakan dan Mengoptimalkan Pipeline Data dengan Polars

Dataset permintaan yang dibersihkan

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"),
    )
)
Menskalakan dan Mengoptimalkan Pipeline Data dengan Polars

Kumpulkan lalu tulis

clean_requests.collect(
    engine="streaming"
).write_parquet(
    "311_requests_clean.parquet"
)
Menskalakan dan Mengoptimalkan Pipeline Data dengan Polars

Sink ke Parquet

clean_requests.sink_parquet(
    "311_requests_clean.parquet",


)
Menskalakan dan Mengoptimalkan Pipeline Data dengan Polars

Mengatur sink

clean_requests.sink_parquet(
    "311_requests_clean.parquet",
    compression="zstd",
    row_group_size=100_000,
)
Menskalakan dan Mengoptimalkan Pipeline Data dengan Polars

Metode sink lain

df.sink_csv("output.csv")
df.sink_ndjson("output.ndjson")
Menskalakan dan Mengoptimalkan Pipeline Data dengan Polars

Output terpartisi

clean_requests.sink_parquet(
    pl.PartitionBy(



    ),

)
Menskalakan dan Mengoptimalkan Pipeline Data dengan Polars

Output terpartisi

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

)
Menskalakan dan Mengoptimalkan Pipeline Data dengan Polars

Output terpartisi

clean_requests.sink_parquet(
    pl.PartitionBy(
        "311_requests_clean/",
        key="CREATED_DATE",
        max_rows_per_file=1_000_000,
    ),
    mkdir=True,
)
Menskalakan dan Mengoptimalkan Pipeline Data dengan Polars

Output terpartisi

311_requests_clean/
  CREATED_DATE=2025-12-31/
    00000000.parquet
  CREATED_DATE=2026-01-01/
    00000000.parquet
Menskalakan dan Mengoptimalkan Pipeline Data dengan Polars

Partisi multi-kolom

clean_requests.sink_parquet(
    pl.PartitionBy(
        "311_requests_clean/",
        key=["STATUS", "CREATED_DATE"],
    ),
    mkdir=True,
)
Menskalakan dan Mengoptimalkan Pipeline Data dengan Polars

Membangun sink lazy

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





Menskalakan dan Mengoptimalkan Pipeline Data dengan Polars

Membangun sink lazy

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,
)
Menskalakan dan Mengoptimalkan Pipeline Data dengan Polars

Multiplexing sink

pl.collect_all(
    [internal_sink, public_sink],
)
Menskalakan dan Mengoptimalkan Pipeline Data dengan Polars

Output yang ditulis

311_requests_internal.parquet
311_requests_public.parquet
Menskalakan dan Mengoptimalkan Pipeline Data dengan Polars

Sink batch kustom

def send_batch(batch: pl.DataFrame) -> None:
    print(batch.height)
    batch_json = batch.write_json()
    # send batch_json to an API
Menskalakan dan Mengoptimalkan Pipeline Data dengan Polars

Sink batch kustom

clean_requests.sink_batches(
    send_batch,
    chunk_size=50_000,
)
50000
50000
...
29268
Menskalakan dan Mengoptimalkan Pipeline Data dengan Polars

Alur kerja sink

  • Gunakan sink_parquet() untuk menulis hasil kueri besar langsung ke disk
  • Gunakan PartitionBy() untuk membagi output menjadi dataset terpartisi
  • Gunakan sink lazy dengan collect_all() untuk merencanakan dan mengeksekusi beberapa output sekaligus
  • Gunakan sink_batches() saat tiap batch dikirim ke fungsi kustom

Tampilan isometrik sistem pengurutan data yang mendistribusikan data ke laci partisi yang rapi

Menskalakan dan Mengoptimalkan Pipeline Data dengan Polars

Ayo berlatih!

Menskalakan dan Mengoptimalkan Pipeline Data dengan Polars

Preparing Video For Download...