Dask DataFrame

Lập trình song song với Dask trong Python

James Fulton

Climate Informatics Researcher

pandas DataFrame so với Dask DataFrame

import pandas as pd

# Đọc một tệp csv pandas_df = pd.read_csv( "dataset/chunk1.csv" )

Thao tác này đọc ngay một tệp CSV.

import dask.dataframe as dd

# Đọc lười tất cả các tệp csv dask_df = dd.read_csv( "dataset/*.csv" )

Thao tác này đọc lười tất cả tệp CSV trong thư mục dataset.

Lập trình song song với Dask trong Python

Dask DataFrame

print(dask_df)
Cấu trúc Dask DataFrame:
              ID     col1     col2    col3     col4     ...
npartitions=3                               
           int64   object   object   int64  float64     ...
             ...      ...      ...     ...      ...     ...
             ...      ...      ...     ...      ...     ...
             ...      ...      ...     ...      ...     ...
Tên Dask: getitem, 3 tác vụ
Lập trình song song với Dask trong Python

Đồ thị tác vụ Dask DataFrame

dask.visualize(dask_df)

Đồ thị tác vụ cho thấy 3 thao tác tải.

Lập trình song song với Dask trong Python

Kiểm soát kích thước khối

# Đặt kích thước bộ nhớ tối đa cho mỗi khối
dask_df = dd.read_csv("dataset/*.csv", blocksize="10MB")

print(dask_df)
Cấu trúc Dask DataFrame:
              ID     col1     col2    col3     col4     ...
npartitions=7                               
           int64   object   object   int64  float64     ...
             ...      ...      ...     ...      ...     ...
             ...      ...      ...     ...      ...     ...
             ...      ...      ...     ...      ...     ...
Tên Dask: getitem, 7 tác vụ
Lập trình song song với Dask trong Python

Giải thích partition

# Đặt kích thước bộ nhớ tối đa cho mỗi khối
dask_df = dd.read_csv("dataset/*.csv", blocksize="10MB")

Vì sao có 7 partition?

size      file
  9M      dataset/chunk1.csv
 18M      dataset/chunk2.csv
 32M      dataset/chunk3.csv
Lập trình song song với Dask trong Python

Giải thích partition

# Đặt kích thước bộ nhớ tối đa cho mỗi khối
dask_df = dd.read_csv("dataset/*.csv", blocksize="10MB")

Vì sao có 7 partition?

size      file
  9M      dataset/chunk1.csv    # thành 1 partition
 18M      dataset/chunk2.csv    # thành 2 partition
 32M      dataset/chunk3.csv    # thành 4 partition
Lập trình song song với Dask trong Python

Phân tích với Dask DataFrame

  • Chọn cột
    col1 = dask_df['col1']
    
  • Gán cột
    dask_df['double_col1'] = 2 * col1
    
  • Phép toán số học, ví dụ
    dask_df.std()
    dask_df.min()
    
  • Groupby
    dask_df.groupby(col1).mean()
    
  • Cả các hàm bạn đã dùng trước đó
    dask_df.nlargest(n=3, columns='col1')
    
Lập trình song song với Dask trong Python

Datetime và các chức năng khác của pandas

import pandas as pd

# Chuyển chuỗi sang định dạng datetime
pd.to_datetime(pandas_df['start_date'])


# Truy cập thuộc tính 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

# Chuyển chuỗi sang định dạng datetime
dd.to_datetime(dask_df['start_date'])


# Truy cập thuộc tính 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
Lập trình song song với Dask trong Python

Chuyển kết quả thành không-lười

# Hiển thị 5 dòng
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     ...
# Chuyển Dask DataFrame (lười) sang pandas DataFrame trong bộ nhớ
results_df = df.compute()
Lập trình song song với Dask trong Python

Ghi kết quả trực tiếp ra tệp

# 7 partition (khối) nên có 7 tệp đầu ra
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
Lập trình song song với Dask trong Python

Định dạng tệp nhanh hơn - Parquet

# Đọc từ parquet
dask_df = dd.read_parquet('dataset_parquet')

# Ghi ra parquet
dask_df.to_parquet('answer_parquet')
  • Parquet đọc nhanh hơn CSV nhiều lần
  • Có thể ghi nhanh hơn
Lập trình song song với Dask trong Python

Ayo berlatih!

Lập trình song song với Dask trong Python

Preparing Video For Download...