Lập trình song song với Dask trong Python
James Fulton
Climate Informatics Researcher
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.
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ụ
dask.visualize(dask_df)

# Đặ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ụ
# Đặ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
# Đặ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
col1 = dask_df['col1']
dask_df['double_col1'] = 2 * col1
dask_df.std()
dask_df.min()
dask_df.groupby(col1).mean()
dask_df.nlargest(n=3, columns='col1')
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
# 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()
# 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
# Đọc từ parquet
dask_df = dd.read_parquet('dataset_parquet')
# Ghi ra parquet
dask_df.to_parquet('answer_parquet')
Lập trình song song với Dask trong Python