Julia molte allocazioni che utilizzano Distributed e SharedArrays con @ sync / @ async
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
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 eseguireaddprocsin qualche parte dello script che esegue l'inizializzazione o forse i processi vengono aggiunti avviando iljuliaprocesso 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