Dask DataFrame入門
pd.read_csvとpd.DataFrameをDaskの同等機能に置き換え、compute()で実行を開始し、タスクグラフをプロファイリングします。
「Dask DataFrame入門」はCoddyKit上の無料Pandas & NumPy Academyレッスンです。 これはレッスン3/4です。 下記で完全なレッスンを無料で読むことができます。その後、ブラウザ内の組み込みコードエディタと24時間対応のAIチューターでハンズオン演習できます。 これはPandas & NumPy Academy学習パスの一部であり、ウェブとCoddyKitアプリ全体で進捗が同期されます。 Pandas & NumPy Academyコースには全4レッスンが含まれています。
Daskとは
Daskは、NumPyやPandasを拡張し、RAMより大きなデータセットを扱えるようにするPython向けの並列計算ライブラリです。そのdask.dataframeモジュールは、Pandasとほぼ同一のDataFrame APIを提供します。ただし、Daskは操作をすぐに実行するのではなくタスクグラフを構築し、.compute()を呼び出したときに遅延実行します。これにより、コードをほとんど変更せずに、複数のコア、さらには複数のマシンに処理を分散できます。
Daskのインストールとインポート
Daskはpip install dask[dataframe]でインストールします。インポートにはimport dask.dataframe as ddを使うのが一般的です。内部では、Dask DataFrameは複数の小さなPandas DataFrameに分割され、それぞれが独立して処理されます。Dask DataFrameに対する操作は遅延タスクグラフを作成しますが、.compute()が呼び出されるまで何も実行されません。計算を記述することと実行することを分離する点が、Daskの重要な考え方です。
import dask.dataframe as dd
# Read a large CSV — returns a Dask DataFrame immediately (no data loaded yet)
ddf = dd.read_csv('large_sales.csv')
print(type(ddf)) # dask.dataframe.core.DataFrame
print(ddf.columns.tolist())
print(ddf.dtypes)DaskとPandasの違い:重要なポイント
Pandasでは、すべての操作が即座に eager 実行されます。一方Daskでは、操作の結果として、遅延されている計算を表す別のDaskオブジェクトが返されます。.compute()を呼び出したときに初めて、Daskが実際にデータを読み込み、タスクグラフを実行します。この遅延実行により、Daskは実行前に処理計画を最適化できます。たとえば、連続するフィルターを融合して、データの複数回読み込みを避けられます。料理にたとえるなら、Daskがレシピを書き、.compute()が料理を作る役割を担います。
import dask.dataframe as dd
ddf = dd.read_csv('sales.csv')
# This does NOT run yet — just builds the task graph
filtered = ddf[ddf['amount'] > 1000]
agg = filtered.groupby('region')['amount'].sum()
print(type(agg)) # dask.dataframe.core.Series
# NOW execute everything
result = agg.compute()
print(result)パーティション:中心となる概念
Dask DataFrameはパーティションに分割されます。各パーティションは通常のPandas DataFrameです。デフォルトでは、dd.read_csvは1ファイルにつき1つのパーティションを作成します(大きなファイルでは128 MBごとに1つ)。この設定はblocksizeで変更できます。ddf.npartitionsを確認すると、パーティション数がわかります。パーティションを増やすと並列性が高まりますが、オーバーヘッドも増えます。減らすとオーバーヘッドは小さくなりますが、並列性が制限されます。通常は数百個程度が適切なバランスです。
import dask.dataframe as dd
ddf = dd.read_csv('data/*.csv') # Read multiple CSV files at once
print('Number of partitions:', ddf.npartitions)
# Access a single partition as a Pandas DataFrame
first_partition = ddf.get_partition(0).compute()
print('Partition 0 shape:', first_partition.shape)Daskで使えるおなじみのPandas操作
一般的なPandas操作の多くはDaskでも同じように使えます。.head()、.tail()、.describe()、ブールインデックス、.groupby()、.merge()、.assign()には、いずれもDask版があります。最大の違いは、結果を実体化するために.compute()を呼び出す必要があることです。Pandasなら数ミリ秒で処理できる操作でも、Daskではタスクグラフのオーバーヘッドにより数秒かかる場合があります。そのため、小規模なデータにはPandasを使い、RAMに収まらないデータにはDaskを使ってください。
import dask.dataframe as dd
ddf = dd.read_csv('transactions.csv')
# Filtering — same syntax as Pandas
high_value = ddf[ddf['amount'] > 500]
# GroupBy aggregation
by_region = high_value.groupby('region')['amount'].mean()
# Execute
result = by_region.compute()
print(result.sort_values(ascending=False))Globパターンで複数ファイルを読み込む
Daskの便利な機能の1つは、globパターンを使って複数のファイルを一度に読み込めることです。dd.read_csv('data/2024-*.csv')は一致するすべてのファイルを読み込み、ファイルごとに1つのパーティションを作成します。これは、月別や日別のファイルに分割して保存されたデータに最適です。データレイクでよく使われる構成です。Daskはスキーマを自動的にそろえるため、Pandasで手動ループと結合を行う場合と同等の処理を、より簡単に実行できます。
import dask.dataframe as dd
# Read all monthly files at once
ddf = dd.read_csv('sales/2024-*.csv',
dtype={'order_id': 'int32',
'amount': 'float32'})
print(f'Partitions: {ddf.npartitions}') # One per file
print(f'Total rows (lazy): {len(ddf)}') # This triggers a compute!タスクグラフのvisualize()メソッド
複雑なDaskパイプラインを実行する前に、result.visualize()を呼び出してタスクグラフを確認できます。これにより、すべての計算手順を示すPNG図が生成されます。Daskが実行する処理を理解したり、予想外の遅さをデバッグしたりする際に役立ちます。グラフには、パーティションがフィルター、groupby、集計の各処理をどのように通過するかが表示されるため、重複した計算も簡単に見つけられます。graphvizパッケージが必要です。
import dask.dataframe as dd
ddf = dd.read_csv('orders.csv')
pipeline = (
ddf[ddf['status'] == 'completed']
.groupby('product_id')['revenue']
.sum()
)
# Visualise the task graph (saves to PNG)
# pipeline.visualize('task_graph.png')
# Check number of tasks in the graph
print('Number of tasks:', len(pipeline.__dask_graph__()))map_partitionsでカスタム関数を適用する
Dask DataFrameにカスタムPandas関数を適用する必要がある場合は、ddf.map_partitions(func)を使います。これは各パーティションに対してfuncを個別に適用し、新しいDask DataFrameを返します。関数には通常のPandas DataFrameが渡され、関数はPandas DataFrameを返す必要があります。これはdf.apply()に相当するDaskの機能であり、Pandas専用のコードをDaskと連携させる方法です。
import dask.dataframe as dd
import pandas as pd
def normalise_chunk(df):
df = df.copy()
df['amount_norm'] = (df['amount'] - df['amount'].mean()) / df['amount'].std()
return df
ddf = dd.read_csv('data.csv')
normalised = ddf.map_partitions(normalise_chunk)
result = normalised[['order_id', 'amount_norm']].compute()
print(result.head())Daskのスケジューラーオプション
Daskには、タスクの実行方法を制御する複数のスケジューラーがあります。'synchronous'スケジューラーは、現在のスレッドでタスクを逐次実行します(デバッグに便利です)。'threads'スケジューラーはスレッドプールを使用します(I/Oバウンドの処理に適しています)。'processes'スケジューラーは、CPUバウンドの処理向けに複数のプロセスを起動します(PythonのGILを回避できます)。Daskの分散クラスターを使うと、複数のマシンで実行できます。スケジューラーはcompute(scheduler='threads')のように指定します。
import dask.dataframe as dd
ddf = dd.read_csv('data.csv')
agg = ddf.groupby('category')['sales'].sum()
# Choose scheduler based on workload
result_sync = agg.compute(scheduler='synchronous') # sequential, easy to debug
result_threads = agg.compute(scheduler='threads') # parallel I/O
result_processes = agg.compute(scheduler='processes') # parallel CPUDaskとPandasの間で変換する
大規模なデータセットをDaskで処理した後、最終的な分析や可視化のために集計結果をPandasに移すことはよくあります。.compute()を使うと、Dask DataFrameをPandasに変換できます。逆方向の変換では、dd.from_pandas(df, npartitions=4)によってPandas DataFrameをDask DataFrameに変換できます。これは、完全なデータセットに拡張する前に、小規模なデータでDaskコードをテストする場合に便利です。
import pandas as pd
import dask.dataframe as dd
# Start with a small Pandas DF for testing
df_small = pd.DataFrame({'a': range(100), 'b': range(100, 200)})
# Convert to Dask for development/testing
ddf = dd.from_pandas(df_small, npartitions=4)
result = ddf.groupby('a')['b'].sum().compute()
print(type(result)) # pandas.Series
print(result.head())DaskとPandasとSQLの使い分け
Daskが常に適切なツールとは限りません。データがRAMに収まる場合(数GB未満)は、オーバーヘッドが少なく、より単純で高速なPandasを使います。データがRAMを超える一方で、Pandasに似た構文と単一マシンでの並列処理が必要な場合は、Daskを使います。データがリレーショナルデータベースに保存されていて、集計をデータベースエンジンに任せられる場合は、SQLまたはデータベースを使います。ペタバイト規模で複数マシンによる分散処理が必要な場合は、SparkまたはBigQueryを使います。
理解度チェック
このレッスンで学んだデータ分析の概念について理解度を確認します。
レッスンのまとめ
このレッスンでは、Dask DataFramesが、使い慣れたAPIを備えたPandasパーティションの遅延評価コレクションであること、.compute()によってタスクグラフの実行が開始されること、そしてmap_partitionsによって任意のカスタムPandas関数をすべてのパーティションに適用できることを学びました。次は、大規模データセットを保存する際にCSVの代わりに使える、高速な列指向形式であるParquetについて学びます。
よくある質問
「Dask DataFrame入門」レッスンは無料ですか?
はい。「Dask DataFrame入門」の完全なテキストはこのウェブで無料で読めます。インタラクティブに演習し(組み込みコードエディタと24時間対応のAIチューター)、Pandas & NumPy Academyコースの残りをアンロックするには、CoddyKit PROにアップグレードしてください。 Pandas & NumPy Academyコースには全4レッスンが含まれています。
「Dask DataFrame入門」で何を学びますか?
pd.read_csvとpd.DataFrameをDaskの同等機能に置き換え、compute()で実行を開始し、タスクグラフをプロファイリングします。 ブラウザで直接実行するハンズオンコードでPandas & NumPy Academyを演習し、24時間対応のAIチューターがレッスンを進める中での質問に答えます。
Pandas & NumPy Academyを始めるのに経験は必要ですか?
事前経験は必要ありません。CoddyKitのPandas & NumPy Academyは初級者から上級者向けに構成されているため、ここから始めるか最初から始めて、自分のペースで進むことができます。 これはレッスン3/4です。
「Dask DataFrame入門」レッスンにはどのくらい時間がかかりますか?
ほとんどのCoddyKitレッスンは約5~10分かかります。各レッスンはコンパクトでインタラクティブなので、着実に進歩し、ウェブとアプリ全体で正確に前回の場所から再開できます。
このPandas & NumPy Academyレッスンでコードを書いて実行できますか?
はい。すべてのPandas & NumPy Academyレッスンに組み込みコードエディタが含まれているため、ブラウザでリアルコードを書いて実行し、即座のAIフィードバックを取得できます。ローカル設定は不要です。
このコースのすべてのレッスン
- chunksizeによるCSVのストリーミング
- チャンク間の逐次集計
- Dask DataFrame入門
- Parquet:高速な列指向ストレージ