Come accelerare il ciclo "for" in una funzione python?
Ho una funzione var. Voglio conoscere il modo migliore possibile per eseguire il ciclo for (per più coordinate: xs e ys) all'interno di questa funzione rapidamente mediante multiprocessing / elaborazione parallela utilizzando tutti i processori, i core e la memoria RAM del sistema.
È possibile utilizzare il Daskmodulo?
pyshedsla documentazione può essere trovata qui .
import numpy as np
from pysheds.grid import Grid
xs = 82.1206, 72.4542, 65.0431, 83.8056, 35.6744
ys = 25.2111, 17.9458, 13.8844, 10.0833, 24.8306
for (x,y) in zip(xs,ys):
grid = Grid.from_raster('E:/data.tif', data_name='map')
grid.catchment(data='map', x=x, y=y, out_name='catch', recursionlimit=1500, xytype='label')
....
....
results
Risposte
Ho provato a fornire un codice riproducibile di seguito utilizzando dask. È possibile aggiungere la parte di elaborazione principale di pyshedso qualsiasi altra funzione in essa contenuta per un'iterazione parallela più rapida dei parametri.
La documentazione del daskmodulo può essere trovata qui .
import dask
from dask import delayed, compute
from dask.distributed import Client, progress
from pysheds.grid import Grid
client = Client(threads_per_worker=2, n_workers=2) #Choose the number of workers and threads per worker over here to deploy for your task.
xs = 82.1206, 72.4542, 65.0431, 83.8056, 35.6744
ys = 25.2111, 17.9458, 13.8844, 10.0833, 24.8306
#Firstly, a function has to be created, where the iteration of the parameters is involved.
def var(x,y):
grid = Grid.from_raster('data.tif', data_name='map')
grid.catchment(data='map', x=x, y=y, out_name='catch', recursionlimit=1500, xytype='label')
...
...
return (result)
#Now calling the function in a 'dask' way.
lazy_results = []
for (x,y) in zip(xs,ys):
lazy_result = dask.delayed(var)(x,y)
lazy_results.append(lazy_result)
#Final command to execute the function var(x,y) and get the result.
dask.compute(*lazy_results)
Non hai pubblicato un collegamento al tuo image1.tiffile, quindi il codice di esempio riportato di seguito utilizza pysheds/data/dem.tifdahttps://github.com/mdbartos/pyshedsL'idea di base è dividere i parametri di input xse , ysnel tuo caso, in sottoinsiemi, quindi assegnare a ciascuna CPU un sottoinsieme diverso su cui lavorare.
main()calcola la soluzione due volte, una in sequenza e una in parallelo, quindi confronta le soluzioni di ciascuna. C'è una certa inefficienza nella soluzione parallela poiché il file immagine verrà letto da ciascuna CPU, quindi c'è spazio per miglioramenti (cioè, leggere il file immagine al di fuori della porzione parallela quindi dare l' gridoggetto risultante a ciascuna istanza).
import numpy as np
from pysheds.grid import Grid
from dask.distributed import Client
from dask import delayed, compute
xs = 10, 20, 30, 40, 50, 60, 70, 80, 90, 100
ys = 25, 35, 45, 55, 65, 75, 85, 95, 105, 115, 125
def var(image_file, x_in, y_in):
grid = Grid.from_raster(image_file, data_name='map')
variable_avg = []
for (x,y) in zip(x_in,y_in):
grid.catchment(data='map', x=x, y=y, out_name='catch')
variable = grid.view('catch', nodata=np.nan)
variable_avg.append( np.array(variable).mean() )
return(variable_avg)
def var_parallel(n_cpu, image_file, x_in, y_in):
tasks = []
for cpu in range(n_cpu):
x_in = xs[cpu::n_cpu] # eg, cpu = 0: x_in = (10, 40, 70, 100)
y_in = ys[cpu::n_cpu] #
tasks.append( delayed(var)(image_file, x_in, y_in) )
ans = compute(tasks)
# reassemble solution in the right order
par_avg = [None]*len(xs)
for cpu in range(n_cpu):
par_avg[cpu::n_cpu] = ans[0][cpu]
print('AVG (parallel) =',par_avg)
return par_avg
def main():
image_file = 'pysheds/data/dem.tif'
# sequential solution:
seq_avg = var(image_file, xs, ys)
print('AVG (sequential)=',seq_avg)
# parallel solution:
n_cpu = 3
dask_client = Client(n_workers=n_cpu)
par_avg = var_parallel(n_cpu, image_file, xs, ys)
dask_client.shutdown()
print('max error=',
max([ abs(seq_avg[i]-par_avg[i]) for i in range(len(seq_avg))]))
if __name__ == '__main__': main()