Julia molte allocazioni che utilizzano Distributed e SharedArrays con @ sync / @ async

Sep 11 2020

Sto cercando di capire come utilizzare il pacchetto Distributed insieme a SharedArrays per eseguire operazioni parallele con julia. A titolo di esempio sto prendendo un semplice metodo medio Montecarlo

using Distributed 
using SharedArrays
using Statistics

const NWorkers = 2
const Ns = Int(1e6)


function parallelRun()

    addprocs(NWorkers)
    procsID = workers()

    A = SharedArray{Float64,1}(Ns)
    println("starting loop")

    for i=1:2:Ns
        #parallel block
        @sync for p=1:NWorkers
                @async A[i+p-1] = remotecall_fetch(rand,procsID[p]);
        end
    end

    println(mean(A))
end


function singleRun()
    A = zeros(Ns)
    for i=1:Ns
        A[i] = rand()
    end
 
    println(mean(A))
end

Tuttavia, se io @time entrambe le funzioni ottengo

julia> @time singleRun()
0.49965531193003165
  0.009762 seconds (17 allocations: 7.630 MiB)
julia> @time parallelRun()
0.4994892300029917
 46.319737 seconds (66.99 M allocations: 2.665 GiB, 1.01% gc time)

In particolare ci sono molte più allocazioni nella versione parallela, il che rende il codice molto più lento.

Mi sto perdendo qualcosa?

A proposito, il motivo per cui sto usando @sync e @async (anche se non sono necessari in questo framework poiché ogni campione può essere calcolato in ordine casuale) è solo perché vorrei applicare la stessa strategia per risolvere numericamente una PDE parabolica con qualcosa sulla linea di

    for t=1:time_steps

        #parallel block
        @sync for p=1:NWorkers
                @async remotecall(make_step_PDE,procsID[p],p);
        end
    end

dove ogni lavoratore indicizzato da p dovrebbe lavorare su un insieme disgiunto di indici della mia equazione.

Grazie in anticipo

Risposte

2 PrzemyslawSzufel Sep 11 2020 at 05:57

Esistono i seguenti problemi nel codice:

  • Stai generando un'attività remota per ogni valore di ia e questo è solo costoso e alla fine ci vuole molto. Fondamentalmente la regola pratica è usare la @distributedmacro per il bilanciamento del carico tra i lavoratori, questo condividerà equamente il lavoro.
  • Non inserire mai la addprocstua funzione di lavoro perché ogni volta che la esegui, ogni volta che aggiungi nuovi processi - anche la generazione di un nuovo processo Julia richiede molto tempo e questo è stato incluso nelle tue misurazioni. In pratica questo significa che vuoi eseguire addprocsin qualche parte dello script che esegue l'inizializzazione o forse i processi vengono aggiunti avviando il juliaprocesso con -po --machine-fileparametro
  • Infine, esegui @timesempre sempre due volte: nella prima misurazione @timevengono misurati anche i tempi di compilazione e la compilazione in un ambiente distribuito richiede molto più tempo rispetto a un singolo processo.

La tua funzione dovrebbe essere più o meno simile a questa

using Distributed, SharedArrays
addprocs(4)
@everywhere using Distributed, SharedArrays
function parallelRun(Ns)
    A = SharedArray{Float64,1}(Ns)
    @sync @distributed for i=1:Ns
         A[i] = rand();
    end
    println(mean(A))
end

Potresti anche considerare di suddividere completamente i dati tra i lavoratori. Questo in alcuni scenari è meno soggetto a bug e ti consente di distribuire su molti nodi:

using Distributed, DistributedArrays
addprocs(4)
@everywhere using Distributed, DistributedArrays
function parallelRun2(Ns)
    d = dzeros(Ns) #creates an array distributed evenly on all workers
    @sync @distributed for i in 1:Ns
        p = localpart(d)
        p[((i-1) % Int(Ns/nworkers())+1] = rand()
    end
    println(mean(d))
end