การ sink คิวรีขนาดใหญ่ลงดิสก์

การปรับขนาดและเพิ่มประสิทธิภาพ Data Pipeline ด้วย Polars

Liam Brannigan

Data Scientist & Polars Contributor

ปัญหาผลลัพธ์ขนาดใหญ่

แผนภาพแสดง lazy query ที่สร้างผลลัพธ์แบบ row-preserving ขนาดใหญ่ซึ่งควรเขียนลงดิสก์โดยตรง

การปรับขนาดและเพิ่มประสิทธิภาพ Data Pipeline ด้วย Polars

ชุดข้อมูล request ที่ทำความสะอาดแล้ว

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

clean_requests = (
    requests





)
การปรับขนาดและเพิ่มประสิทธิภาพ Data Pipeline ด้วย Polars

ชุดข้อมูล request ที่ทำความสะอาดแล้ว

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"),
    )
)
การปรับขนาดและเพิ่มประสิทธิภาพ Data Pipeline ด้วย Polars

Collect แล้วเขียน

clean_requests.collect(
    engine="streaming"
).write_parquet(
    "311_requests_clean.parquet"
)
การปรับขนาดและเพิ่มประสิทธิภาพ Data Pipeline ด้วย Polars

การ sink ลง Parquet

clean_requests.sink_parquet(
    "311_requests_clean.parquet",


)
การปรับขนาดและเพิ่มประสิทธิภาพ Data Pipeline ด้วย Polars

ปรับแต่งการ sink

clean_requests.sink_parquet(
    "311_requests_clean.parquet",
    compression="zstd",
    row_group_size=100_000,
)
การปรับขนาดและเพิ่มประสิทธิภาพ Data Pipeline ด้วย Polars

เมธอด sink อื่น ๆ

df.sink_csv("output.csv")
df.sink_ndjson("output.ndjson")
การปรับขนาดและเพิ่มประสิทธิภาพ Data Pipeline ด้วย Polars

ผลลัพธ์แบบ partition

clean_requests.sink_parquet(
    pl.PartitionBy(



    ),

)
การปรับขนาดและเพิ่มประสิทธิภาพ Data Pipeline ด้วย Polars

ผลลัพธ์แบบ partition

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

)
การปรับขนาดและเพิ่มประสิทธิภาพ Data Pipeline ด้วย Polars

ผลลัพธ์แบบ partition

clean_requests.sink_parquet(
    pl.PartitionBy(
        "311_requests_clean/",
        key="CREATED_DATE",
        max_rows_per_file=1_000_000,
    ),
    mkdir=True,
)
การปรับขนาดและเพิ่มประสิทธิภาพ Data Pipeline ด้วย Polars

ผลลัพธ์แบบ partition

311_requests_clean/
  CREATED_DATE=2025-12-31/
    00000000.parquet
  CREATED_DATE=2026-01-01/
    00000000.parquet
การปรับขนาดและเพิ่มประสิทธิภาพ Data Pipeline ด้วย Polars

การแบ่ง partition หลายคอลัมน์

clean_requests.sink_parquet(
    pl.PartitionBy(
        "311_requests_clean/",
        key=["STATUS", "CREATED_DATE"],
    ),
    mkdir=True,
)
การปรับขนาดและเพิ่มประสิทธิภาพ Data Pipeline ด้วย Polars

การสร้าง lazy sink

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





การปรับขนาดและเพิ่มประสิทธิภาพ Data Pipeline ด้วย Polars

การสร้าง lazy 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,
)
การปรับขนาดและเพิ่มประสิทธิภาพ Data Pipeline ด้วย Polars

การ multiplex sink

pl.collect_all(
    [internal_sink, public_sink],
)
การปรับขนาดและเพิ่มประสิทธิภาพ Data Pipeline ด้วย Polars

ไฟล์ที่เขียนออกมา

311_requests_internal.parquet
311_requests_public.parquet
การปรับขนาดและเพิ่มประสิทธิภาพ Data Pipeline ด้วย Polars

Custom batch sink

def send_batch(batch: pl.DataFrame) -> None:
    print(batch.height)
    batch_json = batch.write_json()
    # send batch_json to an API
การปรับขนาดและเพิ่มประสิทธิภาพ Data Pipeline ด้วย Polars

Custom batch sink

clean_requests.sink_batches(
    send_batch,
    chunk_size=50_000,
)
50000
50000
...
29268
การปรับขนาดและเพิ่มประสิทธิภาพ Data Pipeline ด้วย Polars

ขั้นตอนการใช้ sink

  • ใช้ sink_parquet() เพื่อเขียนผลลัพธ์คิวรีขนาดใหญ่ลงดิสก์โดยตรง
  • ใช้ PartitionBy() เพื่อแบ่งผลลัพธ์เป็น partitioned dataset
  • ใช้ lazy sink ร่วมกับ collect_all() เพื่อวางแผนและรันหลาย output พร้อมกัน
  • ใช้ sink_batches() เมื่อต้องการส่งแต่ละ batch ไปยังฟังก์ชันที่กำหนดเอง

มุมมองแบบ isometric ของระบบจัดเรียงข้อมูลที่กระจายข้อมูลลงใน partition แบบเป็นระเบียบ

การปรับขนาดและเพิ่มประสิทธิภาพ Data Pipeline ด้วย Polars

มาฝึกกันเถอะ!

การปรับขนาดและเพิ่มประสิทธิภาพ Data Pipeline ด้วย Polars

Preparing Video For Download...