Parquetファイルの操作

Polars によるデータパイプラインのスケーリングと最適化

Liam Brannigan

Data Scientist & Polars Contributor

ファイルストレージ形式

DataFrameの図

Polars によるデータパイプラインのスケーリングと最適化

CSV形式

CSVが行単位で格納される様子を示す図。各レコードが1行全体にわたって保存されている。

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の行グループ

Parquetファイルが行グループに分割されている図。

Polars によるデータパイプラインのスケーリングと最適化

行グループを使ったフィルタリング

(
    pl.scan_parquet("311_Service_Requests.parquet")
    .filter(pl.col("CREATED_DATE") < pl.datetime(2020, 1, 1))
)
Polars によるデータパイプラインのスケーリングと最適化

Parquetの行グループ

行グループを持つParquetファイルの図。

Polars によるデータパイプラインのスケーリングと最適化

Parquetの行グループ

行グループと統計情報を持つ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...