Bir python işlevinde 'for' döngüsü nasıl hızlandırılır?

Sep 09 2020

Benim bir işlevim var var. Sistemin sahip olduğu tüm işlemcileri, çekirdekleri ve RAM belleğini kullanarak çoklu işlem / paralel işleme yoluyla bu işlev içinde for döngüsünü (çoklu koordinatlar için: xs ve ys) hızlı bir şekilde çalıştırmanın mümkün olan en iyi yolunu bilmek istiyorum.

DaskModül kullanarak mümkün mü ?

pyshedsdokümantasyon burada bulunabilir .

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

Yanıtlar

1 SaiKiran Sep 16 2020 at 21:54

Kullanarak aşağıda tekrarlanabilir bir kod vermeye çalıştım dask. pyshedsParametrelerin daha hızlı paralel yinelemesi için ana işlem parçasını veya içindeki diğer işlevleri ekleyebilirsiniz .

daskModülün dokümantasyonu burada bulunabilir .

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

image1.tifDosyanıza bir bağlantı göndermediniz, bu nedenle aşağıdaki örnek kod pysheds/data/dem.tif,https://github.com/mdbartos/pyshedsTemel fikir, girdi parametrelerini xsve yssizin durumunuzda alt kümelere ayırmak, ardından her bir CPU'ya üzerinde çalışacak farklı bir alt küme vermektir.

main()çözümü iki kez, bir kez sırayla ve bir kez paralel olarak hesaplar ve ardından her birinden gelen çözümleri karşılaştırır. Paralel çözümde bir miktar verimsizlik vardır, çünkü görüntü dosyası her bir CPU tarafından okunacaktır, bu nedenle iyileştirme için yer vardır (yani, paralel bölümün dışındaki görüntü dosyasını okuyun ve ardından ortaya çıkan gridnesneyi her bir örneğe verin).

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