使用 Parquet 文件

使用 Polars 扩展与优化数据流水线

Liam Brannigan

Data Scientist & Polars Contributor

文件存储格式

数据框示意图

使用 Polars 扩展与优化数据流水线

CSV 格式

示意图:CSV 为逐行存储,每条记录占一整行。

使用 Polars 扩展与优化数据流水线

Parquet 格式

示意图:Parquet 为列式存储,相同列的值分组存放。

使用 Polars 扩展与优化数据流水线

转换归档数据

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

)
使用 Polars 扩展与优化数据流水线

转换归档数据

(
    pl.read_csv("311_Service_Requests.csv", try_parse_dates=True)
    .write_parquet("311_Service_Requests.parquet")
)
使用 Polars 扩展与优化数据流水线

检查 Parquet 文件

pl.read_parquet_schema(
    "311_Service_Requests.parquet"
)
Schema(
    "CREATED_DATE": Datetime,
    "TYPE": String,
    "STATUS": String,
    "DEPARTMENT": String,
    "WARD": Int64,
    ...
)
使用 Polars 扩展与优化数据流水线

CSV vs Parquet

$$

CSV
  • 4.8 GB
  • 14 秒

$$

Parquet
  • 0.6 GB
  • 1.4 秒

csv 与 parquet 对比示意图

使用 Polars 扩展与优化数据流水线

Parquet 不适合追加行

new_request = {"TYPE": "Pothole", "STATUS": "Open", ...}
  • 追加用 CSV
  • 读取用 Parquet
使用 Polars 扩展与优化数据流水线

Parquet 行组(row groups)

示意图:Parquet 文件被划分为行组。

使用 Polars 扩展与优化数据流水线

利用行组过滤

(
    pl.scan_parquet("311_Service_Requests.parquet")
    .filter(pl.col("CREATED_DATE") < pl.datetime(2020, 1, 1))
)
使用 Polars 扩展与优化数据流水线

Parquet 行组(row groups)

带有行组的 Parquet 文件示意图。

使用 Polars 扩展与优化数据流水线

Parquet 行组(row groups)

带有行组与统计信息的 Parquet 文件示意图。

使用 Polars 扩展与优化数据流水线

并行策略扫描

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

    )
)
使用 Polars 扩展与优化数据流水线

并行策略扫描

(
    pl.scan_parquet(
        "311_Service_Requests.parquet",
        parallel="row_groups",
    )
)
  • columns
  • prefiltered
  • none
使用 Polars 扩展与优化数据流水线

控制 Parquet 写入

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

)
  • 文件更小
  • 读写更慢
使用 Polars 扩展与优化数据流水线

控制 Parquet 写入

department_counts.write_parquet(
    "department_counts.parquet",
    compression_level=3,
    row_group_size=25_000,
)
  • 文件更小
  • 读写更慢
  • 需读取更多统计信息
使用 Polars 扩展与优化数据流水线

开始练习!

使用 Polars 扩展与优化数据流水线

Preparing Video For Download...