Использование процессов и потоков

Параллельное программирование с Dask на Python

James Fulton

Climate Informatics Researcher

Планировщик Dask по умолчанию

Потоки

  • Dask arrays
  • Dask DataFrames
  • Отложенные пайплайны, созданные с помощью dask.delayed()

Процессы

  • Dask bags
Параллельное программирование с Dask на Python

Выбор планировщика

# Use default
result = x.compute()

result = dask.compute(x)
# Use threads result = x.compute(scheduler='threads')
result = dask.compute(x, scheduler='threads')
# Use processes result = x.compute(scheduler='processes')
result = dask.compute(x, scheduler='processes')
Параллельное программирование с Dask на Python

Итоги: потоки и процессы

Потоки

  • Запускаются очень быстро
  • Не требуют передачи данных
  • Ограничены GIL: только один поток читает код в каждый момент времени

Процессы

  • Требуют времени на запуск
  • Медленная передача данных
  • Каждый процесс имеет собственный GIL и не ждёт очереди на чтение кода
Параллельное программирование с Dask на Python

Создание локального кластера

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)
Параллельное программирование с Dask на Python

Создание локального кластера

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)
Параллельное программирование с Dask на Python

Простой локальный кластер

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)
Параллельное программирование с Dask на Python

Создание клиента

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>
Параллельное программирование с Dask на Python

Быстрое создание клиента

Создать кластер и передать его в клиент

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, ...>
Параллельное программирование с Dask на Python

Использование кластера

client = Client(processes=True)

# Default uses the client
result = x.compute()

# Can still change to other schedulers result = x.compute(scheduler='threads')
# Can explicitly use client result = client.compute(x)
Параллельное программирование с Dask на Python

Другие типы кластеров

  • LocalCluster() — кластер на вашем компьютере.
  • Кластеры других типов распределяют вычисления между разными машинами
Параллельное программирование с Dask на Python

Давайте потренируемся!

Параллельное программирование с Dask на Python

Preparing Video For Download...