Używanie procesów i wątków

Programowanie równoległe z Dask w Pythonie

James Fulton

Climate Informatics Researcher

Domyślny harmonogram Dask

Wątki

  • Tablice Dask
  • DataFrames Dask
  • Potoki opóźnione z dask.delayed()

Procesy

  • Worki Dask
Programowanie równoległe z Dask w Pythonie

Wybór harmonogramu

# 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')
Programowanie równoległe z Dask w Pythonie

Powtórzenie — wątki i procesy

Wątki

  • Bardzo szybkie uruchamianie
  • Nie wymagają transferu danych
  • Ograniczone przez GIL — tylko jeden wątek może czytać kod naraz

Procesy

  • Wolniejsze uruchamianie
  • Wolny transfer danych
  • Każdy ma własny GIL — nie muszą czekać na dostęp do kodu
Programowanie równoległe z Dask w Pythonie

Tworzenie klastra lokalnego

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)
Programowanie równoległe z Dask w Pythonie

Tworzenie klastra lokalnego

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)
Programowanie równoległe z Dask w Pythonie

Prosty klaster lokalny

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)
Programowanie równoległe z Dask w Pythonie

Tworzenie klienta

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>
Programowanie równoległe z Dask w Pythonie

Uproszczone tworzenie klienta

Utwórz klaster, a następnie przekaż go do klienta

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

client = Client(cluster)

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

Utwórz klienta, który sam utworzy własny klaster

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



print(client)
<Client: ... processes=4 threads=8, ...>
Programowanie równoległe z Dask w Pythonie

Korzystanie z klastra

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)
Programowanie równoległe z Dask w Pythonie

Inne rodzaje klastrów

  • LocalCluster() — klaster na lokalnym komputerze.
  • Inne typy klastrów rozkładają obliczenia na wiele komputerów
Programowanie równoległe z Dask w Pythonie

Czas na ćwiczenia!

Programowanie równoległe z Dask w Pythonie

Preparing Video For Download...