Dask
Dask
Dask는 기존 Python 및 PyData 생태계를 확장하는 병렬·분산 컴퓨팅 라이브러리예요. 특히 Dask DataFrame으로 pandas 워크플로를 확장할 수 있습니다.
출처: 문서
본문
Dask는 기존 Python 및 PyData 생태계를 확장하는 병렬·분산 컴퓨팅 라이브러리입니다.
특히 Dask DataFrame을 사용해 pandas 워크플로를 확장할 수 있어요. Dask DataFrame은 대규모 테이블형 데이터를 다루기 위해 pandas를 병렬화합니다. pandas API를 밀접하게 반영해 단일 데이터셋에서 테스트하다 전체 데이터셋을 처리하는 것으로 간단히 전환할 수 있어요. Dask는 Hugging Face Datasets의 기본 포맷인 Parquet와 특히 효과적입니다. 풍부한 데이터 타입, 효율적인 컬럼형 필터링, 압축을 지원하기 때문입니다.
Dask의 좋은 실용적 사용 사례는 데이터셋에 대한 데이터 처리나 모델 추론을 분산 방식으로 실행하는 것입니다. 예를 들어 Coiled의 훌륭한 블로그 포스트 Scaling AI-Based Data Processing with Hugging Face + Dask를 참고하세요.
읽기와 쓰기 (Read and Write)
Dask는 fsspec으로 원격 데이터를 읽고 쓰므로, Hugging Face 경로(hf://)를 사용해 Hub에서 데이터를 읽고 쓸 수 있어요.
먼저 Hugging Face 계정으로 로그인해야 합니다. 예:
hf auth login
그런 다음 데이터셋 리포지토리를 만들고:
from huggingface_hub import HfApi
HfApi().create_repo(repo_id="username/my_dataset", repo_type="dataset")
마지막으로 Dask에서 Hugging Face 경로를 사용할 수 있어요. Dask DataFrame은 Hugging Face의 Parquet에 대한 분산 쓰기를 지원하며, 커밋으로 데이터셋 변경을 추적합니다:
import dask.dataframe as dd
df.to_parquet("hf://datasets/username/my_dataset")
# or write in separate directories if the dataset has train/validation/test splits
df_train.to_parquet("hf://datasets/username/my_dataset/train")
df_valid.to_parquet("hf://datasets/username/my_dataset/validation")
df_test .to_parquet("hf://datasets/username/my_dataset/test")
이것은 파일당 하나의 커밋을 만들므로 업로드 후 히스토리를 스쿼시(squash)하는 것을 권장합니다:
from huggingface_hub import HfApi
HfApi().super_squash_history(repo_id=repo_id, repo_type="dataset")
이렇게 하면 Dask 데이터셋을 Parquet 포맷으로 담은 username/my_dataset 데이터셋 리포지토리가 생성됩니다. 나중에 다시 불러올 수 있어요:
import dask.dataframe as dd
df = dd.read_parquet("hf://datasets/username/my_dataset")
# or read from separate directories if the dataset has train/validation/test splits
df_train = dd.read_parquet("hf://datasets/username/my_dataset/train")
df_valid = dd.read_parquet("hf://datasets/username/my_dataset/validation")
df_test = dd.read_parquet("hf://datasets/username/my_dataset/test")
Hugging Face 경로와 구현 방식에 대한 자세한 내용은 클라이언트 라이브러리의 HfFileSystem 문서를 참고하세요.
데이터 처리 (Process data)
Dask로 데이터셋을 병렬 처리하려면 먼저 pandas DataFrame 또는 Series용 데이터 처리 함수를 정의한 다음, Dask map_partitions 함수로 이 함수를 데이터셋의 모든 파티션에 병렬로 적용하면 됩니다:
def dummy_count_words(texts):
return pd.Series([len(text.split(" ")) for text in texts])
pandas 문자열 메서드를 사용한 비슷한 함수(더 빠름):
def dummy_count_words(texts):
return texts.str.count(" ")
pandas에서는 텍스트 컬럼에 이 함수를 사용할 수 있어요:
# pandas API
df["num_words"] = dummy_count_words(df.text)
Dask에서는 모든 파티션에서 이 함수를 실행할 수 있습니다:
# Dask API: run the function on every partition
df["num_words"] = df.text.map_partitions(dummy_count_words, meta=int)
함수 출력의 pandas Series 또는 DataFrame 타입인 meta도 제공해야 한다는 점을 알아두세요. Dask DataFrame은 lazy API를 사용하기 때문에 필요합니다. Dask는 .compute()가 호출될 때만 데이터 처리를 실행하므로, 그동안 새 컬럼의 타입을 알기 위해 meta 인자가 필요해요.
Predicate 및 Projection Pushdown
Hugging Face에서 Parquet 데이터를 읽을 때 Dask는 Parquet 파일의 메타데이터를 자동으로 활용해 필요 없는 파일이나 행 그룹을 건너뜁니다. 예를 들어 Parquet 포맷의 Hugging Face 데이터셋에 필터(predicate)를 적용하거나 컬럼 하위 집합(projection)을 선택하면, Dask는 Parquet 파일의 메타데이터를 읽어 필요 없는 부분을 다운로드하지 않고 버립니다.
이는 쿼리 최적화를 지원하는 Dask DataFrame API 재구현 덕분에 가능하며, Dask를 더 빠르고 견고하게 만듭니다.
예를 들어 이 FineWeb-Edu 하위 집합은 많은 Parquet 파일을 포함합니다. 최근 CC 덤프의 텍스트만 유지하도록 데이터셋을 필터링할 수 있다면 Dask는 대부분의 파일을 건너뛰고 필터와 일치하는 데이터만 다운로드합니다:
import dask.dataframe as dd
df = dd.read_parquet("hf://datasets/HuggingFaceFW/fineweb-edu/sample/10BT/*.parquet")
# Dask will skip the files or row groups that don't
# match the query without downloading them.
df = df[df.dump >= "CC-MAIN-2023"]
Dask는 또한 계산에 필요한 컬럼만 읽고 나머지는 건너뜁니다. 예를 들어 코드 후반에 컬럼을 drop하면 필요하지 않다면 파이프라인 초반에 불러오지 않아요. 컬럼 하위 집합을 조작하거나 분석할 때 유용합니다:
# Dask will download the 'dump' and 'token_count' needed
# for the filtering and computation and skip the other columns.
df.token_count.mean().compute()
클라이언트 (Client)
dask의 대부분 기능은 병렬 계산을 실행하기 위한 클러스터 또는 로컬 Client에 최적화되어 있습니다:
import dask.dataframe as dd
from distributed import Client
if __name__ == "__main__": # needed for creating new processes
client = Client()
df = dd.read_parquet(...)
...
로컬 사용에서 Client는 기본적으로 multiprocessing을 쓰는 Dask LocalCluster를 사용합니다. LocalCluster의 multiprocessing을 수동으로 구성할 수 있어요:
from dask.distributed import Client, LocalCluster
cluster = LocalCluster(n_workers=8, threads_per_worker=8)
client = Client(cluster)
Client 없이 로컬에서 기본 스레드 스케줄러를 사용하면 DataFrame이 특정 연산 후 느려질 수 있음을 알아두세요(자세한 내용은 여기).
로컬 또는 클라우드 클러스터 설정에 대한 더 많은 정보는 Deploying Dask 문서에서 확인하세요.
더 알아보기 (Learn more)
hf://datasets/<repo> 경로로 Dask가 Hub의 Parquet 데이터를 분산 읽기·쓰기하고, map_partitions로 병렬 처리하며, meta 인자로 lazy API의 컬럼 타입을 알려줘요. predicate/projection pushdown 덕분에 불필요한 파일·컬럼을 다운로드하지 않습니다. 클러스터 배포와 쿼리 최적화는 Dask 문서를 참고하세요.