Dask 데이터프레임

Python에서 Dask로 병렬 프로그래밍

James Fulton

Climate Informatics Researcher

pandas 데이터프레임 vs. Dask 데이터프레임

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 데이터프레임

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 데이터프레임 태스크 그래프

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 데이터프레임으로 분석

  • 열 선택
    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')
    
Python에서 Dask로 병렬 프로그래밍

Datetime과 기타 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...