プロセスとスレッドの活用

Pythonで学ぶDaskによる並列プログラミング

James Fulton

Climate Informatics Researcher

Dask のデフォルトスケジューラ

スレッド

  • Dask Arrays
  • Dask DataFrames
  • dask.delayed() で作成した遅延パイプライン

プロセス

  • Dask Bags
Pythonで学ぶDaskによる並列プログラミング

スケジューラの選択

# 既定を使用
result = x.compute()

result = dask.compute(x)
# スレッドを使用 result = x.compute(scheduler='threads')
result = dask.compute(x, scheduler='threads')
# プロセスを使用 result = x.compute(scheduler='processes')
result = dask.compute(x, scheduler='processes')
Pythonで学ぶDaskによる並列プログラミング

復習:スレッド vs. プロセス

スレッド

  • 起動が非常に速い
  • データ転送が不要
  • GIL によって、同時に実行できるのは1スレッドのみ

プロセス

  • 立ち上げに時間がかかる
  • データ転送が遅い
  • 各プロセスに独自の GIL があり、同時並行で実行可能
Pythonで学ぶDaskによる並列プログラミング

ローカルクラスターの作成

from dask.distributed import LocalCluster

cluster = LocalCluster(
    processes=True, 
    n_workers=2,
    threads_per_worker=2
)

print(cluster)
LocalCluster(..., workers=2, threads=4, memory=31.38 GiB)
Pythonで学ぶDaskによる並列プログラミング

ローカルクラスターの作成

from dask.distributed import LocalCluster

cluster = LocalCluster(
    processes=False, 
    n_workers=2,
    threads_per_worker=2
)

print(cluster)
LocalCluster(..., workers=2, threads=4, memory=31.38 GiB)
Pythonで学ぶDaskによる並列プログラミング

シンプルなローカルクラスター

cluster = LocalCluster(processes=True)

print(cluster)
LocalCluster(..., workers=4 threads=8, memory=31.38 GiB)
cluster = LocalCluster(processes=False)

print(cluster)
LocalCluster(..., workers=1 threads=8, memory=31.38 GiB)
Pythonで学ぶDaskによる並列プログラミング

クライアントの作成

from dask.distributed import Client, LocalCluster
cluster = LocalCluster(
    processes=True, 
    n_workers=4,
    threads_per_worker=2
)

client = Client(cluster)
print(client)
<Client: 'tcp://127.0.0.1:61391' processes=4 threads=8, memory=31.38 GiB>
Pythonで学ぶDaskによる並列プログラミング

かんたんにクライアントを作成

クラスターを作成し、クライアントに渡す

cluster = LocalCluster(
    processes=True, 
    n_workers=4,
    threads_per_worker=2
)

client = Client(cluster)

print(client)
<Client: ... processes=4 threads=8, ...>

クライアント作成時にクラスターも自動作成

client = Client(
    processes=True, 
    n_workers=4,
    threads_per_worker=2
)



print(client)
<Client: ... processes=4 threads=8, ...>
Pythonで学ぶDaskによる並列プログラミング

クラスターの利用

client = Client(processes=True)

# 既定で client を使用
result = x.compute()

# 別スケジューラに変更可能 result = x.compute(scheduler='threads')
# 明示的に client を使用 result = client.compute(x)
Pythonで学ぶDaskによる並列プログラミング

他のクラスター種別

  • LocalCluster():自分のマシン上のクラスター
  • 他の種類は複数マシンに計算を分散
Pythonで学ぶDaskによる並列プログラミング

Passons à la pratique !

Pythonで学ぶDaskによる並列プログラミング

Preparing Video For Download...