Büyük sorguları diske aktarma

Polars ile Veri Hatlarını Ölçeklendirme ve Optimize Etme

Liam Brannigan

Data Scientist & Polars Contributor

Büyük çıktı sorunu

Tembel bir sorgunun, doğrudan diske yazılması gereken büyük, satırları koruyan bir çıkarım ürettiğini gösteren diyagram.

Polars ile Veri Hatlarını Ölçeklendirme ve Optimize Etme

Temizlenmiş bir istek veri kümesi

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

clean_requests = (
    requests





)
Polars ile Veri Hatlarını Ölçeklendirme ve Optimize Etme

Temizlenmiş bir istek veri kümesi

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"),
    )
)
Polars ile Veri Hatlarını Ölçeklendirme ve Optimize Etme

Topla sonra yaz

clean_requests.collect(
    engine="streaming"
).write_parquet(
    "311_requests_clean.parquet"
)
Polars ile Veri Hatlarını Ölçeklendirme ve Optimize Etme

Parquet'e sink etme

clean_requests.sink_parquet(
    "311_requests_clean.parquet",


)
Polars ile Veri Hatlarını Ölçeklendirme ve Optimize Etme

Sink'i kontrol etme

clean_requests.sink_parquet(
    "311_requests_clean.parquet",
    compression="zstd",
    row_group_size=100_000,
)
Polars ile Veri Hatlarını Ölçeklendirme ve Optimize Etme

Diğer sink yöntemleri

df.sink_csv("output.csv")
df.sink_ndjson("output.ndjson")
Polars ile Veri Hatlarını Ölçeklendirme ve Optimize Etme

Bölümlendirilmiş çıktı

clean_requests.sink_parquet(
    pl.PartitionBy(



    ),

)
Polars ile Veri Hatlarını Ölçeklendirme ve Optimize Etme

Bölümlendirilmiş çıktı

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

)
Polars ile Veri Hatlarını Ölçeklendirme ve Optimize Etme

Bölümlendirilmiş çıktı

clean_requests.sink_parquet(
    pl.PartitionBy(
        "311_requests_clean/",
        key="CREATED_DATE",
        max_rows_per_file=1_000_000,
    ),
    mkdir=True,
)
Polars ile Veri Hatlarını Ölçeklendirme ve Optimize Etme

Bölümlendirilmiş çıktı

311_requests_clean/
  CREATED_DATE=2025-12-31/
    00000000.parquet
  CREATED_DATE=2026-01-01/
    00000000.parquet
Polars ile Veri Hatlarını Ölçeklendirme ve Optimize Etme

Çok sütunlu bölümler

clean_requests.sink_parquet(
    pl.PartitionBy(
        "311_requests_clean/",
        key=["STATUS", "CREATED_DATE"],
    ),
    mkdir=True,
)
Polars ile Veri Hatlarını Ölçeklendirme ve Optimize Etme

Lazy sink'ler oluşturma

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





Polars ile Veri Hatlarını Ölçeklendirme ve Optimize Etme

Lazy sink'ler oluşturma

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,
)
Polars ile Veri Hatlarını Ölçeklendirme ve Optimize Etme

Sink'leri çoklayarak çalıştırma

pl.collect_all(
    [internal_sink, public_sink],
)
Polars ile Veri Hatlarını Ölçeklendirme ve Optimize Etme

Yazılan çıktılar

311_requests_internal.parquet
311_requests_public.parquet
Polars ile Veri Hatlarını Ölçeklendirme ve Optimize Etme

Özel yığın sink'leri

def send_batch(batch: pl.DataFrame) -> None:
    print(batch.height)
    batch_json = batch.write_json()
    # send batch_json to an API
Polars ile Veri Hatlarını Ölçeklendirme ve Optimize Etme

Özel yığın sink'leri

clean_requests.sink_batches(
    send_batch,
    chunk_size=50_000,
)
50000
50000
...
29268
Polars ile Veri Hatlarını Ölçeklendirme ve Optimize Etme

Sink iş akışı

  • Büyük sorgu sonuçlarını doğrudan diske yazmak için sink_parquet() kullanın
  • Çıktıyı bölümlendirmek için PartitionBy() kullanın
  • Birden çok çıktıyı birlikte planlayıp çalıştırmak için collect_all() ile lazy sink'ler kullanın
  • Her yığın özel bir işlevle gönderilecekse sink_batches() kullanın

Verileri düzenli bölümlere dağıtan bir sıralama sisteminin izometrik görünümü

Polars ile Veri Hatlarını Ölçeklendirme ve Optimize Etme

Hadi pratik yapalım!

Polars ile Veri Hatlarını Ölçeklendirme ve Optimize Etme

Preparing Video For Download...