Використання процесів і потоків

Паралельне програмування з Dask у Python

James Fulton

Climate Informatics Researcher

Стандартний диспетчер задач Dask

Потоки

  • Масиви Dask
  • Dask DataFrames
  • Відкладені конвеєри, створені за допомогою dask.delayed()

Процеси

  • «Мішки» Dask (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

Підсумок — потоки vs процеси

Потоки

  • Дуже швидко запускаються
  • Не потрібно передавати їм дані
  • Обмежені 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...