¿Cómo acelerar el ciclo 'for' en una función de Python?

Sep 09 2020

Tengo una función var. Quiero saber la mejor manera posible de ejecutar el bucle for (para múltiples coordenadas: xs e ys) dentro de esta función rápidamente mediante multiprocesamiento / procesamiento paralelo utilizando todos los procesadores, núcleos y memoria RAM que tiene el sistema.

¿Es posible usar el Daskmódulo?

pyshedsLa documentación se puede encontrar aquí .

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

Respuestas

1 SaiKiran Sep 16 2020 at 21:54

Intenté dar un código reproducible a continuación usando dask. Puede agregar la parte de procesamiento principal de la pyshedso cualquier otra función en ella para una iteración paralela más rápida de los parámetros.

La documentación del daskmódulo se puede encontrar aquí .

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)
1 AlDanial Sep 13 2020 at 07:47

No publicó un enlace a su image1.tifarchivo, por lo que el código de muestra a continuación usa pysheds/data/dem.tifdehttps://github.com/mdbartos/pyshedsLa idea básica es dividir los parámetros de entrada xsy , ysen su caso, en subconjuntos, luego darle a cada CPU un subconjunto diferente para trabajar.

main()calcula la solución dos veces, una secuencialmente y otra en paralelo, luego compara las soluciones de cada una. Existe cierta ineficiencia en la solución paralela ya que el archivo de imagen será leído por cada CPU, por lo que hay margen de mejora (es decir, lea el archivo de imagen fuera de la parte paralela y luego proporcione el gridobjeto resultante a cada instancia).

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()