Dask DataFrame

Pythonで学ぶDaskによる並列プログラミング

James Fulton

Climate Informatics Researcher

pandas DataFrame と Dask DataFrame の違い

import pandas as pd

# 単一のCSVを読み込み pandas_df = pd.read_csv( "dataset/chunk1.csv" )

これは単一のCSVを即時に読み込みます。

import dask.dataframe as dd

# 複数CSVを遅延読み込み dask_df = dd.read_csv( "dataset/*.csv" )

これはdatasetフォルダ内のすべてのCSVを遅延読み込みします。

Pythonで学ぶDaskによる並列プログラミング

Dask DataFrame

print(dask_df)
Dask DataFrame 構造:
              ID     col1     col2    col3     col4     ...
npartitions=3                               
           int64   object   object   int64  float64     ...
             ...      ...      ...     ...      ...     ...
             ...      ...      ...     ...      ...     ...
             ...      ...      ...     ...      ...     ...
Dask Name: getitem, 3 tasks
Pythonで学ぶDaskによる並列プログラミング

Dask DataFrame のタスクグラフ

dask.visualize(dask_df)

タスクグラフは3つの読み込み処理を示す。

Pythonで学ぶDaskによる並列プログラミング

ブロックサイズの制御

# チャンクの最大メモリを指定
dask_df = dd.read_csv("dataset/*.csv", blocksize="10MB")

print(dask_df)
Dask DataFrame 構造:
              ID     col1     col2    col3     col4     ...
npartitions=7                               
           int64   object   object   int64  float64     ...
             ...      ...      ...     ...      ...     ...
             ...      ...      ...     ...      ...     ...
             ...      ...      ...     ...      ...     ...
Dask Name: getitem, 7 tasks
Pythonで学ぶDaskによる並列プログラミング

パーティションの説明

# チャンクの最大メモリを指定
dask_df = dd.read_csv("dataset/*.csv", blocksize="10MB")

なぜ7パーティション?

size      file
  9M      dataset/chunk1.csv
 18M      dataset/chunk2.csv
 32M      dataset/chunk3.csv
Pythonで学ぶDaskによる並列プログラミング

パーティションの説明

# チャンクの最大メモリを指定
dask_df = dd.read_csv("dataset/*.csv", blocksize="10MB")

なぜ7パーティション?

size      file
  9M      dataset/chunk1.csv    # 1パーティションになる
 18M      dataset/chunk2.csv    # 2パーティションになる
 32M      dataset/chunk3.csv    # 4パーティションになる
Pythonで学ぶDaskによる並列プログラミング

Dask DataFrameでの分析

  • 列を選択
    col1 = dask_df['col1']
    
  • 列の代入
    dask_df['double_col1'] = 2 * col1
    
  • 数学演算 例:
    dask_df.std()
    dask_df.min()
    
  • groupby
    dask_df.groupby(col1).mean()
    
  • 以前使った関数も可
    dask_df.nlargest(n=3, columns='col1')
    
Pythonで学ぶDaskによる並列プログラミング

日時と他のpandas機能

import pandas as pd

# 文字列をdatetimeに変換
pd.to_datetime(pandas_df['start_date'])


# datetime属性にアクセス pandas_df['start_date'].dt.year pandas_df['start_date'].dt.day pandas_df['start_date'].dt.hour pandas_df['start_date'].dt.minute
import dask.dataframe as dd

# 文字列をdatetimeに変換
dd.to_datetime(dask_df['start_date'])


# datetime属性にアクセス dask_df['start_date'].dt.year dask_df['start_date'].dt.day dask_df['start_date'].dt.hour dask_df['start_date'].dt.minute
Pythonで学ぶDaskによる並列プログラミング

結果を非遅延にする

# 先頭5行を表示
print(dask_df.head())
        ID     double_col1     col1    col2     col3     ...
0   543795              20       10     436        0     ...
1   874535              24       12     268        0     ...
2   781326              62       31     211        0     ...
3   112457              18        9     898        1     ...
4   103256             142       71     663        0     ...
# 遅延のDask DataFrameをメモリ上のpandas DataFrameへ変換
results_df = df.compute()
Pythonで学ぶDaskによる並列プログラミング

結果を直接ファイルへ出力

# 7パーティション(チャンク)なので出力は7ファイル
dask_df.to_csv('answer/part-*.csv')
part-0.csv
part-1.csv
part-2.csv
part-3.csv
part-4.csv
part-5.csv
part-6.csv
Pythonで学ぶDaskによる並列プログラミング

より高速なファイル形式 - Parquet

# Parquetから読み込み
dask_df = dd.read_parquet('dataset_parquet')

# Parquetへ保存
dask_df.to_parquet('answer_parquet')
  • ParquetはCSVより読み込みが数倍高速
  • 書き込みも高速な場合あり
Pythonで学ぶDaskによる並列プログラミング

Ayo berlatih!

Pythonで学ぶDaskによる並列プログラミング

Preparing Video For Download...