Utiliser des processus et des threads

Programmation parallèle avec Dask en Python

James Fulton

Climate Informatics Researcher

Planificateur par défaut de Dask

Threads

  • Tableaux Dask
  • DataFrames Dask
  • Chaînes différées créées avec dask.delayed()

Processus

  • Sacs Dask (bags)
Programmation parallèle avec Dask en Python

Choisir le planificateur

# Utiliser la valeur par défaut
result = x.compute()

result = dask.compute(x)
# Utiliser des threads result = x.compute(scheduler='threads')
result = dask.compute(x, scheduler='threads')
# Utiliser des processus result = x.compute(scheduler='processes')
result = dask.compute(x, scheduler='processes')
Programmation parallèle avec Dask en Python

Récapitulatif — threads vs processus

Threads

  • Très rapides à lancer
  • Aucun transfert de données nécessaire
  • Limités par le GIL, qui n'autorise qu'un thread à lire le code à la fois

Processus

  • Plus longs à initialiser
  • Transfert de données plus lent
  • Chacun a son propre GIL ; pas besoin d'alterner la lecture du code
Programmation parallèle avec Dask en Python

Créer un amas local

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)
Programmation parallèle avec Dask en Python

Créer un amas local

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)
Programmation parallèle avec Dask en Python

Amas local simple

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)
Programmation parallèle avec Dask en Python

Créer un client

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>
Programmation parallèle avec Dask en Python

Créer un client facilement

Créer l'amas, puis le passer au client

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

client = Client(cluster)

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

Créer un client qui créera son propre amas

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



print(client)
<Client: ... processes=4 threads=8, ...>
Programmation parallèle avec Dask en Python

Utiliser l'amas

client = Client(processes=True)

# Par défaut, utilise le client
result = x.compute()

# On peut quand même changer de planificateur result = x.compute(scheduler='threads')
# Utiliser explicitement le client result = client.compute(x)
Programmation parallèle avec Dask en Python

Autres types d'amas

  • LocalCluster() — Un amas sur votre ordinateur.
  • D'autres types d'amas répartissent le calcul sur plusieurs ordinateurs
Programmation parallèle avec Dask en Python

Passons à la pratique !

Programmation parallèle avec Dask en Python

Preparing Video For Download...