Đẩy truy vấn lớn xuống đĩa

Mở rộng và tối ưu hóa pipeline dữ liệu với Polars

Liam Brannigan

Data Scientist & Polars Contributor

Bài toán đầu ra lớn

Sơ đồ cho thấy một truy vấn lười tạo ra bản trích lớn giữ nguyên số hàng và nên ghi trực tiếp xuống đĩa.

Mở rộng và tối ưu hóa pipeline dữ liệu với Polars

Bộ dữ liệu yêu cầu đã làm sạch

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

clean_requests = (
    requests





)
Mở rộng và tối ưu hóa pipeline dữ liệu với Polars

Bộ dữ liệu yêu cầu đã làm sạch

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"),
    )
)
Mở rộng và tối ưu hóa pipeline dữ liệu với Polars

Collect rồi ghi

clean_requests.collect(
    engine="streaming"
).write_parquet(
    "311_requests_clean.parquet"
)
Mở rộng và tối ưu hóa pipeline dữ liệu với Polars

Sink sang Parquet

clean_requests.sink_parquet(
    "311_requests_clean.parquet",


)
Mở rộng và tối ưu hóa pipeline dữ liệu với Polars

Điều khiển sink

clean_requests.sink_parquet(
    "311_requests_clean.parquet",
    compression="zstd",
    row_group_size=100_000,
)
Mở rộng và tối ưu hóa pipeline dữ liệu với Polars

Các phương thức sink khác

df.sink_csv("output.csv")
df.sink_ndjson("output.ndjson")
Mở rộng và tối ưu hóa pipeline dữ liệu với Polars

Xuất phân vùng

clean_requests.sink_parquet(
    pl.PartitionBy(



    ),

)
Mở rộng và tối ưu hóa pipeline dữ liệu với Polars

Xuất phân vùng

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

)
Mở rộng và tối ưu hóa pipeline dữ liệu với Polars

Xuất phân vùng

clean_requests.sink_parquet(
    pl.PartitionBy(
        "311_requests_clean/",
        key="CREATED_DATE",
        max_rows_per_file=1_000_000,
    ),
    mkdir=True,
)
Mở rộng và tối ưu hóa pipeline dữ liệu với Polars

Xuất phân vùng

311_requests_clean/
  CREATED_DATE=2025-12-31/
    00000000.parquet
  CREATED_DATE=2026-01-01/
    00000000.parquet
Mở rộng và tối ưu hóa pipeline dữ liệu với Polars

Phân vùng đa cột

clean_requests.sink_parquet(
    pl.PartitionBy(
        "311_requests_clean/",
        key=["STATUS", "CREATED_DATE"],
    ),
    mkdir=True,
)
Mở rộng và tối ưu hóa pipeline dữ liệu với Polars

Tạo sink lười

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





Mở rộng và tối ưu hóa pipeline dữ liệu với Polars

Tạo sink lười

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,
)
Mở rộng và tối ưu hóa pipeline dữ liệu với Polars

Ghép nhiều sink

pl.collect_all(
    [internal_sink, public_sink],
)
Mở rộng và tối ưu hóa pipeline dữ liệu với Polars

Đầu ra đã ghi

311_requests_internal.parquet
311_requests_public.parquet
Mở rộng và tối ưu hóa pipeline dữ liệu với Polars

Sink theo lô tùy chỉnh

def send_batch(batch: pl.DataFrame) -> None:
    print(batch.height)
    batch_json = batch.write_json()
    # send batch_json to an API
Mở rộng và tối ưu hóa pipeline dữ liệu với Polars

Sink theo lô tùy chỉnh

clean_requests.sink_batches(
    send_batch,
    chunk_size=50_000,
)
50000
50000
...
29268
Mở rộng và tối ưu hóa pipeline dữ liệu với Polars

Quy trình sink

  • Dùng sink_parquet() để ghi kết quả truy vấn lớn trực tiếp xuống đĩa
  • Dùng PartitionBy() để chia nhỏ đầu ra thành tập dữ liệu phân vùng
  • Dùng sink lười với collect_all() để lập kế hoạch và thực thi nhiều đầu ra cùng lúc
  • Dùng sink_batches() khi mỗi lô cần gửi vào hàm tùy chỉnh

Góc nhìn đẳng trắc của hệ thống sắp xếp dữ liệu phân phối dữ liệu vào các ngăn phân vùng có tổ chức

Mở rộng và tối ưu hóa pipeline dữ liệu với Polars

Ayo berlatih!

Mở rộng và tối ưu hóa pipeline dữ liệu với Polars

Preparing Video For Download...