Làm việc với tệp Parquet

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

Liam Brannigan

Data Scientist & Polars Contributor

Định dạng lưu trữ tệp

Sơ đồ DataFrame

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

Định dạng CSV

Sơ đồ minh họa CSV lưu theo từng hàng, mỗi bản ghi trên một hàng đầy đủ.

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

Định dạng Parquet

Sơ đồ minh họa Parquet lưu theo cột, các giá trị cùng cột được gom lại.

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

Chuyển đổi kho lưu trữ

(
    pl.read_csv("311_Service_Requests.csv", try_parse_dates=True)

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

Chuyển đổi kho lưu trữ

(
    pl.read_csv("311_Service_Requests.csv", try_parse_dates=True)
    .write_parquet("311_Service_Requests.parquet")
)
Mở rộng và tối ưu hóa pipeline dữ liệu với Polars

Kiểm tra tệp Parquet

pl.read_parquet_schema(
    "311_Service_Requests.parquet"
)
Schema(
    "CREATED_DATE": Datetime,
    "TYPE": String,
    "STATUS": String,
    "DEPARTMENT": String,
    "WARD": Int64,
    ...
)
Mở rộng và tối ưu hóa pipeline dữ liệu với Polars

CSV vs Parquet

$$

CSV
  • 4.8 GB
  • 14 giây

$$

Parquet
  • 0.6 GB
  • 1.4 giây

minh họa csv vs parquet

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

Parquet không dùng để thêm hàng

new_request = {"TYPE": "Pothole", "STATUS": "Open", ...}
  • CSV để thêm dòng
  • Parquet để đọc
Mở rộng và tối ưu hóa pipeline dữ liệu với Polars

Nhóm hàng trong Parquet

Sơ đồ cho thấy tệp Parquet chia theo nhóm hàng.

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

Lọc bằng nhóm hàng

(
    pl.scan_parquet("311_Service_Requests.parquet")
    .filter(pl.col("CREATED_DATE") < pl.datetime(2020, 1, 1))
)
Mở rộng và tối ưu hóa pipeline dữ liệu với Polars

Nhóm hàng trong Parquet

Sơ đồ tệp Parquet có nhóm hàng.

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

Nhóm hàng trong Parquet

Sơ đồ tệp Parquet có nhóm hàng và thống kê.

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

Quét với chiến lược song song

(
    pl.scan_parquet(
        "311_Service_Requests.parquet",

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

Quét với chiến lược song song

(
    pl.scan_parquet(
        "311_Service_Requests.parquet",
        parallel="row_groups",
    )
)
  • columns
  • prefiltered
  • none
Mở rộng và tối ưu hóa pipeline dữ liệu với Polars

Kiểm soát ghi Parquet

department_counts.write_parquet(
    "department_counts.parquet",
    compression_level=3,

)
  • Kích thước tệp nhỏ hơn
  • Thời gian đọc/ghi dài hơn
Mở rộng và tối ưu hóa pipeline dữ liệu với Polars

Kiểm soát ghi Parquet

department_counts.write_parquet(
    "department_counts.parquet",
    compression_level=3,
    row_group_size=25_000,
)
  • Kích thước tệp nhỏ hơn
  • Thời gian đọc/ghi dài hơn
  • Nhiều thống kê để đọc hơn
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...